diff --git a/docs/explanation/git-family-object-storage.md b/docs/explanation/git-family-object-storage.md index 95e1502..3fa4146 100644 --- a/docs/explanation/git-family-object-storage.md +++ b/docs/explanation/git-family-object-storage.md @@ -271,12 +271,20 @@ removed deliberately. ### Registry and maintenance worker -`.grasp/staging/.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/.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 diff --git a/src/git/handlers.rs b/src/git/handlers.rs index b396c04..952ba97 100644 --- a/src/git/handlers.rs +++ b/src/git/handlers.rs @@ -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:#}"))) diff --git a/src/git/staging.rs b/src/git/staging.rs index ee08eb7..75cc795 100644 --- a/src/git/staging.rs +++ b/src/git/staging.rs @@ -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, + /// Unsigned ref intentions persisted before receive-pack, independently of + /// the periodic purgatory checkpoint. Wall-clock deadlines survive restarts. + #[serde(default)] + pending: Vec, +} + +#[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 { + 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 { let Some(view) = View::resolve(view)? else { return Ok(false); }; @@ -325,7 +342,25 @@ pub fn stage(view: &Path, owed: &[Tip]) -> Result { 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 { 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"); diff --git a/src/git/staging/recovery.rs b/src/git/staging/recovery.rs new file mode 100644 index 0000000..ddf2def --- /dev/null +++ b/src/git/staging/recovery.rs @@ -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::>()?; + 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); + } +} diff --git a/src/purgatory/mod.rs b/src/purgatory/mod.rs index 983d8c6..19a9525 100644 --- a/src/purgatory/mod.rs +++ b/src/purgatory/mod.rs @@ -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. diff --git a/src/server.rs b/src/server.rs index 06edb44..ae7b028 100644 --- a/src/server.rs +++ b/src/server.rs @@ -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 diff --git a/tests/pending_upload_staging.rs b/tests/pending_upload_staging.rs index 463c056..9eaaa79 100644 --- a/tests/pending_upload_staging.rs +++ b/tests/pending_upload_staging.rs @@ -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; +}