diff --git a/docs/explanation/git-family-object-storage.md b/docs/explanation/git-family-object-storage.md index 3fa4146..8389bbb 100644 --- a/docs/explanation/git-family-object-storage.md +++ b/docs/explanation/git-family-object-storage.md @@ -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 placeholder. Legacy records without deadlines receive one grace period that 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 checkpoint. A view whose remaining staged objects belong to pending refs is parked until the next request rather than polled. diff --git a/src/git/staging/recovery.rs b/src/git/staging/recovery.rs index ddf2def..62d921c 100644 --- a/src/git/staging/recovery.rs +++ b/src/git/staging/recovery.rs @@ -1,4 +1,5 @@ //! Restore unsigned-upload expiry independently of the purgatory checkpoint. +use std::path::Path; use std::time::SystemTime; use anyhow::{ensure, Context, Result}; @@ -23,75 +24,96 @@ pub async fn recover_placeholders( Err(error) => return Err(error.into()), }; 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") { 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()? { + if let Err(error) = recover_record(&path, storage, purgatory, database).await { + tracing::error!(record = %path.display(), error = %format!("{error:#}"), + "Cannot recover staged upload expiry; leaving record and data for inspection"); + } + } + Ok(()) +} + +// Failure is local to this record: without a recovered placeholder, its live +// refs remain protected by compaction. Do not delete or rename suspect metadata. +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::>()?; + 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 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 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, ); - let parts: Vec<_> = record - .view - .iter() - .map(|part| part.to_str().context("non-UTF-8 view path")) - .collect::>()?; - 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(()) } @@ -124,6 +146,48 @@ mod tests { (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] async fn uncheckpointed_uploads_recover_original_expiry_and_are_reclaimed() { for prs in [false, true] { diff --git a/tests/pending_upload_staging.rs b/tests/pending_upload_staging.rs index 9eaaa79..ac13021 100644 --- a/tests/pending_upload_staging.rs +++ b/tests/pending_upload_staging.rs @@ -403,3 +403,31 @@ async fn crash_without_purgatory_checkpoint_still_expires_unsigned_upload() { assert_eq!(family_objects(&served.family), family_before); 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(®istry) + .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; +}