fix(git): isolate staged expiry recovery failures

A malformed or unreadable staging record should not keep the entire relay offline. Skipping that record creates no new expiry authority and preserves live refs and their history.

Recover each record independently, log failures with the record path, and continue recovering healthy views. Leave suspect metadata and data intact for inspection. Registry enumeration failure remains fatal; this change does not add automatic metadata repair or batch retained-root writes.

Validation: reproduced the corrupt JSON failure before the fix. All 18 staging unit tests and 62 integration/helper tests passed, including invalid-path isolation, healthy-record recovery, data preservation and relay restart with corrupt metadata. 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:43:55 +00:00
parent c8fee2ce6a
commit 8c2ee32ee6
3 changed files with 159 additions and 63 deletions
@@ -282,6 +282,10 @@ refs and their durable deadlines. Accepted events and signed events already in
purgatory are preserved. Failed pushes with no matching ref create no purgatory are preserved. Failed pushes with no matching ref create no
placeholder. Legacy records without deadlines receive one grace period that placeholder. Legacy records without deadlines receive one grace period that
is persisted before recovery proceeds, so repeated restarts do not extend it. is persisted before recovery proceeds, so repeated restarts do not extend it.
Unreadable or invalid individual records are logged and skipped without
removing their metadata, refs or objects; healthy records still recover. A
skipped upload with no checkpoint placeholder remains retained for operator
inspection. Failure to enumerate the registry itself still aborts startup.
This closes the window between receive-pack and the periodic purgatory This closes the window between receive-pack and the periodic purgatory
checkpoint. A view whose remaining staged objects belong to pending refs is checkpoint. A view whose remaining staged objects belong to pending refs is
parked until the next request rather than polled. parked until the next request rather than polled.
+127 -63
View File
@@ -1,4 +1,5 @@
//! Restore unsigned-upload expiry independently of the purgatory checkpoint. //! Restore unsigned-upload expiry independently of the purgatory checkpoint.
use std::path::Path;
use std::time::SystemTime; use std::time::SystemTime;
use anyhow::{ensure, Context, Result}; use anyhow::{ensure, Context, Result};
@@ -23,75 +24,96 @@ pub async fn recover_placeholders(
Err(error) => return Err(error.into()), Err(error) => return Err(error.into()),
}; };
for entry in entries { for entry in entries {
let path = entry?.path(); let path = match entry {
Ok(entry) => entry.path(),
Err(error) => {
tracing::error!(%error, "Cannot read staging registry entry; skipping expiry recovery for it");
continue;
}
};
if path.extension().is_none_or(|ext| ext != "json") { if path.extension().is_none_or(|ext| ext != "json") {
continue; continue;
} }
let mut record: Record = serde_json::from_slice(&std::fs::read(&path)?)?; if let Err(error) = recover_record(&path, storage, purgatory, database).await {
ensure!( tracing::error!(record = %path.display(), error = %format!("{error:#}"),
!record.view.as_os_str().is_empty() "Cannot recover staged upload expiry; leaving record and data for inspection");
&& record }
.view }
.components() Ok(())
.all(|part| matches!(part, std::path::Component::Normal(_))), }
"invalid staging view path"
); // Failure is local to this record: without a recovered placeholder, its live
let view_path = storage.git_data_path().join(&record.view); // refs remain protected by compaction. Do not delete or rename suspect metadata.
if !view_path.try_exists()? { async fn recover_record(
path: &Path,
storage: &LocalGitStorage,
purgatory: &Purgatory,
database: &SharedDatabase,
) -> Result<()> {
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()? {
return Ok(());
}
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; continue;
} }
let view = View::resolve(&view_path)?.context("staged view has no family")?; let tip = Tip::new(&reference, &oid);
ensure!( let expires_at =
view.record == std::fs::canonicalize(&path)?, if let Some(pending) = record.pending.iter().find(|pending| pending.tip == tip) {
"staging registry path does not match view" 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,
); );
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(()) Ok(())
} }
@@ -124,6 +146,48 @@ mod tests {
(fixture, owner, oid) (fixture, owner, oid)
} }
#[tokio::test]
async fn bad_records_do_not_block_healthy_recovery_or_delete_data() {
for corrupt in [true, false] {
let (fixture, _, _) = fixture(false);
let bad_path = fixture
.storage
.git_data_path()
.join("unexpected/nested/owner/repo.git");
std::fs::create_dir_all(bad_path.parent().unwrap()).unwrap();
fixture
.storage
.create_thin_view(&fixture.key, &bad_path)
.unwrap();
stage(&bad_path, &[]).unwrap();
let oid = commit(&bad_path, "bad record upload", None);
let bad_ref = super::super::tests::ABANDONED;
stage_upload(&bad_path, &[], &[Tip::new(bad_ref, &oid)]).unwrap();
reference(&bad_path, bad_ref, &oid);
let record_path = View::resolve(&bad_path).unwrap().unwrap().record;
if corrupt {
std::fs::write(&record_path, b"{broken json").unwrap();
}
let before = std::fs::read(&record_path).unwrap();
let purgatory = Purgatory::new(fixture.storage.git_data_path());
recover_placeholders(&fixture.storage, &purgatory, &database())
.await
.unwrap();
assert!(purgatory
.find_pr(PENDING.trim_start_matches("refs/nostr/"))
.is_some());
assert!(purgatory
.find_pr(bad_ref.trim_start_matches("refs/nostr/"))
.is_none());
assert_eq!(
crate::git::get_ref_commit(&bad_path, bad_ref),
Some(oid.clone())
);
assert!(object_exists(&bad_path, &oid).unwrap());
assert_eq!(std::fs::read(&record_path).unwrap(), before);
}
}
#[tokio::test] #[tokio::test]
async fn uncheckpointed_uploads_recover_original_expiry_and_are_reclaimed() { async fn uncheckpointed_uploads_recover_original_expiry_and_are_reclaimed() {
for prs in [false, true] { for prs in [false, true] {
+28
View File
@@ -403,3 +403,31 @@ async fn crash_without_purgatory_checkpoint_still_expires_unsigned_upload() {
assert_eq!(family_objects(&served.family), family_before); assert_eq!(family_objects(&served.family), family_before);
relay.stop().await; relay.stop().await;
} }
#[tokio::test]
async fn corrupt_staging_record_does_not_prevent_relay_restart() {
let mut served = Served::start("corrupt-staging-record").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, "pending with damaged metadata");
served.push(local.path(), &tip, &pending);
let registry = served._persistent.path().join("git/.grasp/staging");
let records: Vec<_> = std::fs::read_dir(&registry)
.unwrap()
.map(|entry| entry.unwrap().path())
.collect();
assert_eq!(records.len(), 1);
let damaged = b"{broken json";
std::fs::write(&records[0], damaged).unwrap();
// restart() kills the process and waits for the new relay to be ready.
served.relay = served.relay.restart().await;
assert_eq!(std::fs::read(&records[0]).unwrap(), damaged);
assert_eq!(
ngit_grasp::git::get_ref_commit(&served.view, &format!("refs/nostr/{}", pending.id)),
Some(tip.clone())
);
assert!(holds_history(&served.view, &tip));
assert!(!holds_history(&served.family, &tip));
served.relay.stop().await;
}