fix(git): recover unsigned upload expiry after crashes

A crash between receive-pack and the purgatory checkpoint left a live staged ref with no expiry owner. Maintenance preserved that ref indefinitely, defeating reclamation of abandoned uploads.

Persist unsigned ref intentions and absolute deadlines in the staging registry before Git runs. Startup reconstructs placeholders from matching live refs before cleanup and request handling, preserving signed events and newer checkpoint deadlines. Legacy records receive one persisted grace period. Compaction drops stale journal entries while retaining owed history.

Recovery assumes exclusive startup access and an available event database; unreadable recovery metadata fails startup rather than authorizing deletion. This change does not add quotas or alter normal placeholder expiry policy.

Validation: 1,179 tests passed across the library, pending_upload_staging, purgatory and grasp06_pr_hosting suites. Regressions cover missing checkpoints, repeated restart deadlines, both endpoint scopes, signed-event preservation, failed pushes and legacy recovery. Workspace all-target clippy with warnings denied, cargo fmt --check and git diff --check passed.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-09-29 11:25:37 +00:00
parent 6de1d47316
commit c8fee2ce6a
7 changed files with 409 additions and 7 deletions
+12 -4
View File
@@ -271,12 +271,20 @@ removed deliberately.
### Registry and maintenance worker
`.grasp/staging/<digest>.json` records each staged view and its owed tips. It
is written and fsynced before receive-pack can write an object. One worker
`.grasp/staging/<digest>.json` records each staged view, its owed tips, and
unsigned ref intentions with absolute expiry deadlines. It is written and
fsynced before receive-pack can write an object. One worker
processes maintenance requests, which come from staged pushes, promotions and
ref deletions, including placeholder expiry. At startup it reads the registry,
so interrupted work resumes. A view whose remaining staged objects belong to
pending refs is parked until the next request rather than polled.
so interrupted work resumes. Before cleanup, sync and request handling start,
recovery reconstructs missing purgatory placeholders from surviving unsigned
refs and their durable deadlines. Accepted events and signed events already in
purgatory are preserved. Failed pushes with no matching ref create no
placeholder. Legacy records without deadlines receive one grace period that
is persisted before recovery proceeds, so repeated restarts do not extend it.
This closes the window between receive-pack and the periodic purgatory
checkpoint. A view whose remaining staged objects belong to pending refs is
parked until the next request rather than polled.
### Limits
+6 -1
View File
@@ -1186,8 +1186,13 @@ pub(crate) async fn stage_push(
})
.map(|(_, new_oid, name)| super::staging::Tip::new(name, new_oid))
.collect();
let unsigned: Vec<_> = pushed_refs
.iter()
.filter(|(_, _, name)| unsigned_refs.contains(name))
.map(|(_, oid, name)| super::staging::Tip::new(name, oid))
.collect();
let view = repo_path.to_owned();
tokio::task::spawn_blocking(move || super::staging::stage(&view, &owed))
tokio::task::spawn_blocking(move || super::staging::stage_upload(&view, &owed, &unsigned))
.await
.map_err(|error| GitError::Storage(error.to_string()))?
.map_err(|error| GitError::Storage(format!("cannot stage push: {error:#}")))
+51 -1
View File
@@ -27,6 +27,8 @@ use serde::{Deserialize, Serialize};
use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat};
mod recovery;
pub use recovery::recover_placeholders;
mod worker;
pub use worker::{request_maintenance, run_worker};
@@ -56,6 +58,16 @@ struct Record {
view: PathBuf,
#[serde(default)]
owed: Vec<Tip>,
/// Unsigned ref intentions persisted before receive-pack, independently of
/// the periodic purgatory checkpoint. Wall-clock deadlines survive restarts.
#[serde(default)]
pending: Vec<PendingUpload>,
}
#[derive(Debug, Serialize, Deserialize)]
struct PendingUpload {
tip: Tip,
expires_at: std::time::SystemTime,
}
/// Outcome of one compaction.
@@ -302,6 +314,11 @@ pub fn owes_history(view: &Path) -> bool {
/// Returns `false` for a legacy repository without a family, which keeps its
/// own objects and is never staged.
pub fn stage(view: &Path, owed: &[Tip]) -> Result<bool> {
stage_upload(view, owed, &[])
}
/// Persist unsigned ref intentions and their expiry before receiving objects.
pub fn stage_upload(view: &Path, owed: &[Tip], unsigned: &[Tip]) -> Result<bool> {
let Some(view) = View::resolve(view)? else {
return Ok(false);
};
@@ -325,7 +342,25 @@ pub fn stage(view: &Path, owed: &[Tip]) -> Result<bool> {
record.owed.push(tip.clone());
}
}
if before != record.owed.len() || !view.record.is_file() {
for tip in unsigned {
ensure!(valid_oid(&tip.oid), "invalid pending object ID");
let id = tip
.reference
.strip_prefix("refs/nostr/")
.context("invalid pending ref")?;
ensure!(
nostr_sdk::prelude::EventId::from_hex(id).is_ok(),
"invalid pending event ID"
);
record
.pending
.retain(|pending| pending.tip.reference != tip.reference);
record.pending.push(PendingUpload {
tip: tip.clone(),
expires_at: std::time::SystemTime::now() + crate::purgatory::DEFAULT_EXPIRY,
});
}
if before != record.owed.len() || !unsigned.is_empty() || !view.record.is_file() {
view.save(&record)?;
}
Ok(true)
@@ -503,6 +538,21 @@ pub fn compact(view: &Path) -> Result<Compaction> {
if !view.record.is_file() {
return Ok(Compaction::Done);
}
// Bound the journal to surviving refs. Failed or expired uploads must not
// accumulate metadata while a sibling keeps the view staged indefinitely.
if let Some(mut record) = view.load()? {
let refs: std::collections::HashMap<_, _> = super::list_refs(&view.path)
.map_err(anyhow::Error::msg)?
.into_iter()
.collect();
let before = record.pending.len();
record
.pending
.retain(|pending| refs.get(&pending.tip.reference) == Some(&pending.tip.oid));
if record.pending.len() != before {
view.save(&record)?;
}
}
let owed = settle_resolved(&view)?;
if owed > 0 {
bail!("{owed} signed tips are still owed to the family");
+259
View File
@@ -0,0 +1,259 @@
//! Restore unsigned-upload expiry independently of the purgatory checkpoint.
use std::time::SystemTime;
use anyhow::{ensure, Context, Result};
use nostr_sdk::prelude::{EventId, FromBech32, PublicKey};
use super::{registry_dir, PendingUpload, Record, Tip, View};
use crate::git::storage::LocalGitStorage;
use crate::nostr::SharedDatabase;
use crate::purgatory::Purgatory;
/// Called during startup before background writers and request handling.
/// Only surviving refs without a known signed event become placeholders.
/// Legacy records get one persisted grace period, never a fresh one per boot.
pub async fn recover_placeholders(
storage: &LocalGitStorage,
purgatory: &Purgatory,
database: &SharedDatabase,
) -> Result<()> {
let entries = match std::fs::read_dir(registry_dir(storage)) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(error) => return Err(error.into()),
};
for entry in entries {
let path = entry?.path();
if path.extension().is_none_or(|ext| ext != "json") {
continue;
}
let mut record: Record = serde_json::from_slice(&std::fs::read(&path)?)?;
ensure!(
!record.view.as_os_str().is_empty()
&& record
.view
.components()
.all(|part| matches!(part, std::path::Component::Normal(_))),
"invalid staging view path"
);
let view_path = storage.git_data_path().join(&record.view);
if !view_path.try_exists()? {
continue;
}
let view = View::resolve(&view_path)?.context("staged view has no family")?;
ensure!(
view.record == std::fs::canonicalize(&path)?,
"staging registry path does not match view"
);
let parts: Vec<_> = record
.view
.iter()
.map(|part| part.to_str().context("non-UTF-8 view path"))
.collect::<Result<_>>()?;
let (owner, prs) = match parts.as_slice() {
[owner, _] => (PublicKey::from_bech32(owner)?, false),
["prs", owner, _] => (PublicKey::from_hex(owner)?, true),
_ => anyhow::bail!("unsupported staged view path"),
};
let refs = crate::git::list_refs(&view.path).map_err(anyhow::Error::msg)?;
for (reference, oid) in refs {
let Some(id) = reference.strip_prefix("refs/nostr/") else {
continue;
};
let Ok(event_id) = EventId::from_hex(id) else {
continue;
};
if database.event_by_id(&event_id).await?.is_some()
|| purgatory
.find_pr(id)
.is_some_and(|entry| entry.event.is_some())
{
continue;
}
let tip = Tip::new(&reference, &oid);
let expires_at =
if let Some(pending) = record.pending.iter().find(|pending| pending.tip == tip) {
pending.expires_at
} else {
// Older registry versions recorded only the view and owed tips.
let expires_at = SystemTime::now() + crate::purgatory::DEFAULT_EXPIRY;
record
.pending
.retain(|pending| pending.tip.reference != reference);
record.pending.push(PendingUpload { tip, expires_at });
view.save(&record)?;
expires_at
};
purgatory.recover_upload_placeholder(
id.to_owned(),
oid,
(owner, view.key.identifier.clone()),
prs,
expires_at,
);
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::super::tests::{commit, reference, Fixture, PENDING};
use super::super::{compact, object_exists, stage, stage_upload};
use super::*;
use nostr_sdk::prelude::*;
use std::sync::Arc;
use std::time::{Duration, Instant};
fn database() -> SharedDatabase {
Arc::new(nostr_memory::MemoryDatabase::unbounded())
}
fn fixture(prs: bool) -> (Fixture, PublicKey, String) {
let owner = Keys::generate().public_key();
let parent = if prs {
format!("prs/{}", owner.to_hex())
} else {
owner.to_bech32().unwrap()
};
let fixture = Fixture::with_view(&parent);
stage(&fixture.view, &[]).unwrap();
let oid = commit(&fixture.view, "pending", None);
stage_upload(&fixture.view, &[], &[Tip::new(PENDING, &oid)]).unwrap();
reference(&fixture.view, PENDING, &oid);
(fixture, owner, oid)
}
#[tokio::test]
async fn uncheckpointed_uploads_recover_original_expiry_and_are_reclaimed() {
for prs in [false, true] {
let (fixture, owner, oid) = fixture(prs);
let view = View::resolve(&fixture.view).unwrap().unwrap();
let mut record = view.load().unwrap().unwrap();
// Inject an elapsed deadline rather than waiting thirty minutes.
record.pending[0].expires_at = SystemTime::UNIX_EPOCH;
view.save(&record).unwrap();
let db = database();
// Repeated recovery must not refresh the durable deadline.
for _ in 0..2 {
let purgatory = Purgatory::new(fixture.storage.git_data_path());
recover_placeholders(&fixture.storage, &purgatory, &db)
.await
.unwrap();
let entry = purgatory
.find_pr(PENDING.trim_start_matches("refs/nostr/"))
.unwrap();
assert!(entry.expires_at <= Instant::now());
assert_eq!(entry.commit, oid);
if prs {
assert_eq!(entry.prs_scope.unwrap().submitter, owner);
} else {
assert_eq!(entry.standard_refs[0].owner, owner);
}
}
let purgatory = Purgatory::new(fixture.storage.git_data_path());
purgatory.set_prs_cleanup_ctx(crate::purgatory::PrsCleanupCtx {
git_data_path: fixture.storage.git_data_path().to_owned(),
repo_init_locks: Default::default(),
});
recover_placeholders(&fixture.storage, &purgatory, &db)
.await
.unwrap();
purgatory
.cleanup_standard_pr_refs(
&crate::nostr::lifecycle::RepositoryLifecycle::in_memory(),
&db,
)
.await;
purgatory.cleanup();
if prs {
assert!(!fixture.view.exists());
} else {
compact(&fixture.view).unwrap();
assert!(!object_exists(&fixture.view, &oid).unwrap());
}
}
}
#[tokio::test]
async fn legacy_registry_gets_only_one_persisted_grace_period() {
let (fixture, _, _) = fixture(false);
let view = View::resolve(&fixture.view).unwrap().unwrap();
let mut record = view.load().unwrap().unwrap();
record.pending.clear();
view.save(&record).unwrap();
let db = database();
let purgatory = Purgatory::new(fixture.storage.git_data_path());
recover_placeholders(&fixture.storage, &purgatory, &db)
.await
.unwrap();
let deadline = view.load().unwrap().unwrap().pending[0].expires_at;
assert!(deadline > SystemTime::now());
let restarted = Purgatory::new(fixture.storage.git_data_path());
recover_placeholders(&fixture.storage, &restarted, &db)
.await
.unwrap();
assert_eq!(
view.load().unwrap().unwrap().pending[0].expires_at,
deadline
);
}
#[tokio::test]
async fn recovery_preserves_signed_events_and_ignores_failed_ref_updates() {
let (fixture, _, oid) = fixture(false);
let keys = Keys::generate();
let event = EventBuilder::new(Kind::GitPullRequest, "signed")
.tags([Tag::custom("c", [&oid])])
.finalize(&keys)
.unwrap();
let reference_name = format!("refs/nostr/{}", event.id);
stage_upload(&fixture.view, &[], &[Tip::new(&reference_name, &oid)]).unwrap();
reference(&fixture.view, &reference_name, &oid);
let db = database();
let purgatory = Purgatory::new(fixture.storage.git_data_path());
purgatory.add_pr(event.clone(), event.id.to_hex(), oid.clone(), false);
recover_placeholders(&fixture.storage, &purgatory, &db)
.await
.unwrap();
assert_eq!(
purgatory.find_pr(&event.id.to_hex()).unwrap().event,
Some(event.clone())
);
db.save_event(&event).await.unwrap();
let restarted = Purgatory::new(fixture.storage.git_data_path());
recover_placeholders(&fixture.storage, &restarted, &db)
.await
.unwrap();
assert!(restarted.find_pr(&event.id.to_hex()).is_none());
assert!(object_exists(&fixture.view, &oid).unwrap());
// Intent was persisted but receive-pack never installed this ref.
let failed = format!("refs/nostr/{}", "ab".repeat(32));
stage_upload(&fixture.view, &[], &[Tip::new(&failed, &oid)]).unwrap();
recover_placeholders(&fixture.storage, &restarted, &db)
.await
.unwrap();
assert!(restarted
.find_pr(failed.trim_start_matches("refs/nostr/"))
.is_none());
}
#[tokio::test]
async fn recovery_preserves_a_newer_checkpoint_deadline() {
let (fixture, owner, oid) = fixture(false);
let purgatory = Purgatory::new(fixture.storage.git_data_path());
let id = PENDING.trim_start_matches("refs/nostr/");
purgatory.recover_upload_placeholder(
id.into(),
oid,
(owner, "repo".into()),
false,
SystemTime::now() + Duration::from_secs(3600),
);
let deadline = purgatory.find_pr(id).unwrap().expires_at;
recover_placeholders(&fixture.storage, &purgatory, &database())
.await
.unwrap();
assert_eq!(purgatory.find_pr(id).unwrap().expires_at, deadline);
}
}
+31 -1
View File
@@ -40,7 +40,7 @@ use std::time::{Duration, Instant, SystemTime};
pub use sync::SyncQueueEntry;
/// Default expiry duration for purgatory entries (30 minutes)
const DEFAULT_EXPIRY: Duration = Duration::from_secs(1800);
pub(crate) const DEFAULT_EXPIRY: Duration = Duration::from_secs(1800);
/// Extended expiry for soft-expired announcements (24 hours).
///
@@ -566,6 +566,36 @@ impl Purgatory {
entry.expires_at = now + DEFAULT_EXPIRY;
}
/// Restore a staged upload before serving requests. Never replace a signed
/// purgatory event or shorten a newer checkpoint's placeholder lifetime.
pub(crate) fn recover_upload_placeholder(
&self,
event_id: String,
commit: String,
destination: (PublicKey, String),
prs: bool,
expires_at: SystemTime,
) {
let (owner, identifier) = destination;
let existing = self.find_pr(&event_id);
if existing.as_ref().is_some_and(|entry| entry.event.is_some()) {
return;
}
let deadline = Instant::now()
+ expires_at
.duration_since(SystemTime::now())
.unwrap_or_default();
let deadline = existing.map_or(deadline, |entry| deadline.max(entry.expires_at));
if prs {
self.add_prs_pr_placeholder(event_id.clone(), commit, owner, identifier);
} else {
self.add_standard_pr_placeholder(event_id.clone(), commit, owner, identifier);
}
if let Some(mut entry) = self.pr_events.get_mut(&event_id) {
entry.expires_at = deadline;
}
}
/// Expire normal-endpoint placeholders under the same lifecycle locks as pushes.
/// Recheck the entry after waiting: an event or a new push may have arrived.
/// Failed Git operations retain the record for the next sweep.
+9
View File
@@ -249,6 +249,15 @@ impl RelayServer {
.await
.context("failed deletion lifecycle startup reconciliation")?;
// Recover upload expiry before cleanup, integrity, sync or HTTP can run.
git::staging::recover_placeholders(
&git::storage::LocalGitStorage::new(config.effective_git_data_path()),
&purgatory,
&relay_runtime.stores.database,
)
.await
.context("failed to recover staged upload expiry")?;
// Make the operator identity discoverable on this relay without
// overwriting identity events the owner already published — locally
// or on the configured user-index relays. Runs after deletion startup
+41
View File
@@ -362,3 +362,44 @@ async fn arriving_state_promotes_an_existing_unsigned_upload() {
async fn waiting_state_promotes_an_unsigned_upload_before_release() {
state_adopts_unsigned_upload(true).await;
}
#[tokio::test]
async fn crash_without_purgatory_checkpoint_still_expires_unsigned_upload() {
let served = Served::start("crash-expiry").await;
let local = tempfile::tempdir().unwrap();
let tip = common::create_test_repo_with_commit(local.path(), CommitVariant::PrTest).unwrap();
let pending = served.pr_event(&tip, "never published");
served.push(local.path(), &tip, &pending);
let family_before = family_objects(&served.family);
let git_data = served._persistent.path().join("git");
let relay_data = served._persistent.path().join("relay");
served.relay.stop().await;
// Reproduce the crash window deterministically: no checkpoint knows this
// upload, and its durable deadline has elapsed while the relay was down.
let checkpoint = git_data.join("purgatory-state.json");
if checkpoint.exists() {
std::fs::remove_file(checkpoint).unwrap();
}
for entry in std::fs::read_dir(git_data.join(".grasp/staging")).unwrap() {
let path = entry.unwrap().path();
let mut record: serde_json::Value =
serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
for upload in record["pending"].as_array_mut().unwrap() {
upload["expires_at"] = serde_json::to_value(std::time::SystemTime::UNIX_EPOCH).unwrap();
}
std::fs::write(path, serde_json::to_vec(&record).unwrap()).unwrap();
}
let relay = TestRelay::start_with_existing_lmdb_paths(git_data, relay_data, None, false).await;
common::wait_for(
"recovered expired upload to be reclaimed",
DEADLINE,
|| async { !ngit_grasp::git::oid_exists(&served.view, &tip) },
)
.await;
assert!(
ngit_grasp::git::get_ref_commit(&served.view, &format!("refs/nostr/{}", pending.id))
.is_none()
);
assert_eq!(family_objects(&served.family), family_before);
relay.stop().await;
}