diff --git a/src/git/sync.rs b/src/git/sync.rs index b96312e..f832945 100644 --- a/src/git/sync.rs +++ b/src/git/sync.rs @@ -1446,8 +1446,10 @@ async fn process_purgatory_announcements( } }; - if let (Some(wp), Some(entry)) = (write_policy, purgatory.find_announcement(&owner, identifier)) - { + if let (Some(wp), Some(entry)) = ( + write_policy, + purgatory.find_announcement(&owner, identifier), + ) { wp.maybe_recover_deleted_repository(&entry.event, identifier) .await; } diff --git a/src/main.rs b/src/main.rs index 30168c6..34984c3 100644 --- a/src/main.rs +++ b/src/main.rs @@ -355,6 +355,7 @@ async fn run_relay(config: Config) -> Result<()> { PathBuf::from(config.effective_git_data_path()), Some(config.domain.clone()), Some(relay_with_db.relay.clone()), + Some(relay_with_db.write_policy.clone()), git_naughty_list.clone(), )); diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index 3cf2dbd..a05b46b 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -457,7 +457,7 @@ impl Nip34WritePolicy { // Process state alignment asynchronously match self .state_policy - .process_state_event(event, is_synced) + .process_state_event(event, is_synced, Some(self)) .await { Ok(policy_result) => { @@ -674,7 +674,7 @@ impl Nip34WritePolicy { // Re-evaluate authorization with the new announcement match self .state_policy - .process_state_event(&entry.event, false) + .process_state_event(&entry.event, false, Some(self)) .await { Ok(WritePolicyResult::Accept) => { diff --git a/src/nostr/policy/state.rs b/src/nostr/policy/state.rs index 1afadb5..59cfe4f 100644 --- a/src/nostr/policy/state.rs +++ b/src/nostr/policy/state.rs @@ -12,6 +12,7 @@ use nostr_relay_builder::prelude::Event; use super::{accepted_purgatory, duplicate, reject_invalid, reject_restricted, PolicyContext}; use crate::git; use crate::git::authorization::fetch_repository_data_with_purgatory; +use crate::nostr::builder::Nip34WritePolicy; use crate::nostr::events::{validate_state, RepositoryAnnouncement, RepositoryState}; /// Result of state policy evaluation @@ -53,6 +54,7 @@ impl StatePolicy { &self, event: &Event, is_synced: bool, + recovery_policy: Option<&Nip34WritePolicy>, ) -> Result { // Parse state to get HEAD and branch info let state = @@ -235,7 +237,7 @@ impl StatePolicy { local_relay.as_ref(), &self.ctx.purgatory, &self.ctx.git_data_path, - None, + recovery_policy, None, ) .await diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index dad1a4f..3bc02f8 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -192,6 +192,7 @@ use std::process::Command; use std::sync::Arc; use tracing::debug; +use crate::nostr::builder::Nip34WritePolicy; use crate::nostr::builder::SharedDatabase; use crate::nostr::events::RepositoryState; use crate::purgatory::Purgatory; @@ -222,6 +223,9 @@ pub struct RealSyncContext { /// Local relay for notifying WebSocket subscribers local_relay: Option, + /// Write policy used for promotion-time recovery hooks. + write_policy: Option, + /// Naughty list tracker for git remote domains with persistent errors git_naughty_list: Arc, } @@ -235,6 +239,7 @@ impl RealSyncContext { /// * `git_data_path` - Base path for git repositories /// * `our_domain` - Our domain to exclude from clone URLs /// * `local_relay` - Local relay for WebSocket notifications + /// * `write_policy` - Write policy for promotion-time recovery hooks /// * `git_naughty_list` - Naughty list tracker for git remote domains pub fn new( purgatory: Arc, @@ -242,6 +247,7 @@ impl RealSyncContext { git_data_path: PathBuf, our_domain: Option, local_relay: Option, + write_policy: Option, git_naughty_list: Arc, ) -> Self { Self { @@ -250,6 +256,7 @@ impl RealSyncContext { git_data_path, our_domain_value: our_domain, local_relay, + write_policy, git_naughty_list, } } @@ -508,8 +515,9 @@ impl SyncContext for RealSyncContext { new_oids: &HashSet, ) -> Result { // Delegate to the unified function from git::sync. - // Pass None for write_policy and rejected_events_index: the purgatory sync path - // already handles hot-cache re-processing via SyncManager::process_event_static. + // Pass configured write_policy for promotion-time recovery hooks. + // Keep rejected_events_index as None: the purgatory sync path already handles + // hot-cache re-processing via SyncManager::process_event_static. let result = crate::git::sync::process_newly_available_git_data( source_repo_path, new_oids, @@ -517,7 +525,7 @@ impl SyncContext for RealSyncContext { self.local_relay.as_ref(), &self.purgatory, &self.git_data_path, - None, + self.write_policy.as_ref(), None, ) .await?; diff --git a/tests/nip09_recovery.rs b/tests/nip09_recovery.rs index dd0ea00..514587b 100644 --- a/tests/nip09_recovery.rs +++ b/tests/nip09_recovery.rs @@ -3,9 +3,9 @@ mod common; use common::{ - announcement_coordinate, build_deletion, publish_served_announcement_for_identifier, - publish_served_repo, publish_served_repo_with_maintainers, publish_served_repo_with_state_event, - create_pr_event, create_state_event, create_test_repo_with_commit, push_to_relay, + announcement_coordinate, build_deletion, create_pr_event, create_state_event, + create_test_repo_with_commit, publish_served_announcement_for_identifier, publish_served_repo, + publish_served_repo_with_maintainers, publish_served_repo_with_state_event, push_to_relay, CommitVariant, TestRelay, }; use grasp_audit::{AuditClient, AuditConfig}; @@ -35,27 +35,123 @@ fn metadata_tag_value(event: &Event, key: &str) -> Option { } fn build_reannouncement(client: &AuditClient, relay: &TestRelay, repo_id: &str) -> Event { + build_reannouncement_with_clone( + client, + relay, + repo_id, + &format!( + "http://{}/{}/{}.git", + relay + .url() + .trim_start_matches("ws://") + .trim_start_matches("wss://"), + client.public_key().to_bech32().expect("pubkey to bech32"), + repo_id + ), + ) +} + +fn build_reannouncement_with_clone( + client: &AuditClient, + relay: &TestRelay, + repo_id: &str, + clone_url: &str, +) -> Event { let relay_url = relay.url().to_string(); - let relay_domain = relay - .url() - .trim_start_matches("ws://") - .trim_start_matches("wss://") - .to_string(); - let npub = client.public_key().to_bech32().expect("pubkey to bech32"); client .event_builder(Kind::GitRepoAnnouncement, "reannounce") .tag(Tag::identifier(repo_id)) .tag(Tag::custom("name", vec![repo_id.to_string()])) - .tag(Tag::custom( - "clone", - vec![format!("http://{}/{}/{}.git", relay_domain, npub, repo_id)], - )) + .tag(Tag::custom("clone", vec![clone_url.to_string()])) .tag(Tag::custom("relays", vec![relay_url])) .build(client.keys()) .expect("build re-announcement") } +async fn publish_served_repo_for_identifier( + client: &AuditClient, + relay: &TestRelay, + repo_id: &str, +) -> String { + let temp_dir = tempfile::tempdir().expect("create temp repo for served repo publish"); + let commit_hash = create_test_repo_with_commit(temp_dir.path(), CommitVariant::StateTest) + .expect("create deterministic commit for served repo"); + + let relay_url = relay.url().to_string(); + let relay_domain = relay_url + .trim_start_matches("ws://") + .trim_start_matches("wss://") + .to_string(); + let npub = client + .public_key() + .to_bech32() + .expect("npub for served repo publish"); + let clone_url = format!("http://{}/{}/{}.git", relay_domain, npub, repo_id); + + let announcement = client + .event_builder(Kind::GitRepoAnnouncement, "") + .tag(Tag::identifier(repo_id)) + .tag(Tag::custom("name", vec![repo_id.to_string()])) + .tag(Tag::custom("clone", vec![clone_url.clone()])) + .tag(Tag::custom("relays", vec![relay_url.clone()])) + .build(client.keys()) + .expect("build served announcement"); + client + .send_event(announcement.clone()) + .await + .expect("send served announcement"); + + let state_event = create_state_event( + client.keys(), + repo_id, + &[("main", commit_hash.as_str())], + &[], + &[clone_url.as_str()], + &[relay_url.as_str()], + ) + .expect("build served state event"); + client + .send_event_and_note_purgatory(state_event.clone()) + .await + .expect("send served state event"); + + push_to_relay(temp_dir.path(), &relay_domain, &npub, repo_id).expect("push served repo data"); + tokio::time::sleep(Duration::from_millis(900)).await; + + assert!( + client + .is_event_on_relay(announcement.id) + .await + .expect("query served announcement"), + "served announcement should be promoted" + ); + assert!( + client + .is_event_on_relay(state_event.id) + .await + .expect("query served state"), + "served state event should be promoted" + ); + + commit_hash +} + +async fn wait_for_event_presence( + client: &AuditClient, + event_id: EventId, + timeout: Duration, +) -> bool { + let deadline = std::time::Instant::now() + timeout; + while std::time::Instant::now() < deadline { + if client.is_event_on_relay(event_id).await.unwrap_or(false) { + return true; + } + tokio::time::sleep(Duration::from_millis(250)).await; + } + false +} + fn run_git_in_repo(repo_path: &Path, args: &[&str]) -> std::process::Output { Command::new("git") .args(args) @@ -146,7 +242,11 @@ async fn setup_parked_announcement_with_failed_ingest_recovery( .await .expect("send pr update"); - set_bare_ref(&repo_path, &format!("refs/nostr/{}", pr.id.to_hex()), &main_commit); + set_bare_ref( + &repo_path, + &format!("refs/nostr/{}", pr.id.to_hex()), + &main_commit, + ); set_bare_ref( &repo_path, &format!("refs/nostr/{}", pr_update.id.to_hex()), @@ -276,11 +376,9 @@ async fn submit_new_state_and_push_for_promotion( #[tokio::test] async fn purgatory_promotion_restores_events_and_pr_git_archive_data() { - let fixture = setup_parked_announcement_with_failed_ingest_recovery( - "promotion-recovery-success", - false, - ) - .await; + let fixture = + setup_parked_announcement_with_failed_ingest_recovery("promotion-recovery-success", false) + .await; std::fs::write(&fixture.archive_path, &fixture.archive_backup) .expect("restore archive bytes before promotion"); @@ -330,11 +428,9 @@ async fn purgatory_promotion_restores_events_and_pr_git_archive_data() { #[tokio::test] async fn purgatory_promotion_recovery_respects_state_tombstones() { - let fixture = setup_parked_announcement_with_failed_ingest_recovery( - "promotion-recovery-tombstone", - true, - ) - .await; + let fixture = + setup_parked_announcement_with_failed_ingest_recovery("promotion-recovery-tombstone", true) + .await; std::fs::write(&fixture.archive_path, &fixture.archive_backup) .expect("restore archive bytes before tombstone promotion"); @@ -378,7 +474,8 @@ async fn repeated_purgatory_promotion_attempts_remain_idempotent() { let first_state = submit_new_state_and_push_for_promotion(&fixture, "idempotent-first").await; // Trigger another post-promotion write flow; recovery should remain a no-op. - let second_reannouncement = build_reannouncement(&fixture.client, &fixture.relay, &fixture.repo_id); + let second_reannouncement = + build_reannouncement(&fixture.client, &fixture.relay, &fixture.repo_id); fixture .client .send_event(second_reannouncement.clone()) @@ -426,13 +523,12 @@ async fn repeated_purgatory_promotion_attempts_remain_idempotent() { #[tokio::test] async fn promotion_with_missing_or_invalid_recovery_artifacts_is_deterministic() { - let fixture = setup_parked_announcement_with_failed_ingest_recovery( - "promotion-recovery-fallback", - false, - ) - .await; + let fixture = + setup_parked_announcement_with_failed_ingest_recovery("promotion-recovery-fallback", false) + .await; - let promoted_state = submit_new_state_and_push_for_promotion(&fixture, "promotion-fallback").await; + let promoted_state = + submit_new_state_and_push_for_promotion(&fixture, "promotion-fallback").await; assert!(fixture .client @@ -476,6 +572,166 @@ async fn promotion_with_missing_or_invalid_recovery_artifacts_is_deterministic() fixture.relay.stop().await; } +#[tokio::test] +async fn state_event_received_route_runs_promotion_recovery_with_tombstones() { + let fixture = setup_parked_announcement_with_failed_ingest_recovery( + "promotion-recovery-state-route", + true, + ) + .await; + + std::fs::write(&fixture.archive_path, &fixture.archive_backup) + .expect("restore archive bytes before state-route promotion"); + + let collaborator = AuditClient::new(fixture.relay.url(), AuditConfig::isolated()) + .await + .expect("create collaborator client"); + + let commit_hash = + publish_served_repo_for_identifier(&collaborator, &fixture.relay, &fixture.repo_id).await; + + let relay_url = fixture.relay.url().to_string(); + let relay_domain = relay_url + .trim_start_matches("ws://") + .trim_start_matches("wss://") + .to_string(); + let collaborator_npub = collaborator + .public_key() + .to_bech32() + .expect("collaborator npub for state-route test"); + let collaborator_clone = format!( + "http://{}/{}/{}.git", + relay_domain, collaborator_npub, fixture.repo_id + ); + + let recovery_state = create_state_event( + fixture.client.keys(), + &fixture.repo_id, + &[("main", commit_hash.as_str())], + &[], + &[collaborator_clone.as_str()], + &[relay_url.as_str()], + ) + .expect("build state-route recovery trigger state event"); + + fixture + .client + .send_event(recovery_state.clone()) + .await + .expect("send state-route recovery trigger state event"); + + assert!( + wait_for_event_presence( + &fixture.client, + fixture.reannouncement.id, + Duration::from_secs(8) + ) + .await, + "state-event route should promote parked re-announcement" + ); + assert!( + wait_for_event_presence(&fixture.client, recovery_state.id, Duration::from_secs(8)).await, + "state-event route should accept recovery trigger state" + ); + assert!( + wait_for_event_presence(&fixture.client, fixture.issue.id, Duration::from_secs(8)).await, + "state-event route should run recovery and restore issue" + ); + assert!( + !fixture + .client + .is_event_on_relay(fixture.old_state_event.id) + .await + .expect("query tombstoned state after state-route recovery"), + "tombstoned state must not resurrect on state-event route" + ); + + fixture.relay.stop().await; +} + +#[tokio::test] +async fn state_event_received_route_fallback_is_deterministic_when_recovery_artifacts_missing() { + let fixture = setup_parked_announcement_with_failed_ingest_recovery( + "promotion-recovery-state-route-fallback", + false, + ) + .await; + + let collaborator = AuditClient::new(fixture.relay.url(), AuditConfig::isolated()) + .await + .expect("create collaborator client for fallback path"); + let commit_hash = + publish_served_repo_for_identifier(&collaborator, &fixture.relay, &fixture.repo_id).await; + + let relay_url = fixture.relay.url().to_string(); + let relay_domain = relay_url + .trim_start_matches("ws://") + .trim_start_matches("wss://") + .to_string(); + let collaborator_npub = collaborator + .public_key() + .to_bech32() + .expect("collaborator npub for fallback state-route test"); + let collaborator_clone = format!( + "http://{}/{}/{}.git", + relay_domain, collaborator_npub, fixture.repo_id + ); + + let fallback_state = create_state_event( + fixture.client.keys(), + &fixture.repo_id, + &[("main", commit_hash.as_str())], + &[], + &[collaborator_clone.as_str()], + &[relay_url.as_str()], + ) + .expect("build fallback state-route trigger state event"); + + fixture + .client + .send_event(fallback_state.clone()) + .await + .expect("send fallback state-route trigger state event"); + + assert!( + wait_for_event_presence( + &fixture.client, + fixture.reannouncement.id, + Duration::from_secs(8) + ) + .await, + "state-event route should promote parked re-announcement in fallback path" + ); + assert!( + wait_for_event_presence(&fixture.client, fallback_state.id, Duration::from_secs(8)).await, + "state-event route should accept fallback trigger state" + ); + assert!( + !fixture + .client + .is_event_on_relay(fixture.issue.id) + .await + .expect("query issue after state-route fallback"), + "state-event route fallback must avoid partial event resurrection" + ); + + let pending = open_holding(&fixture.relay) + .await + .eligible_recovery_records( + &fixture.client.public_key().to_hex(), + &fixture.repo_id, + Timestamp::now(), + DEFAULT_RETENTION, + ) + .await; + assert!( + !pending.is_empty(), + "state-event route fallback should leave deterministic holding records for retries" + ); + + fixture.relay.stop().await; +} + #[tokio::test] async fn reannouncement_restores_events_and_git_archive_happy_path() { let relay = TestRelay::start_with_lmdb().await;