feat(recovery): run promotion recovery on all state promotion routes

This commit is contained in:
DanConwayDev
2026-06-18 09:34:49 +00:00
parent 8206a1057f
commit a52711f195
6 changed files with 308 additions and 39 deletions
+4 -2
View File
@@ -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;
}
+1
View File
@@ -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(),
));
+2 -2
View File
@@ -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) => {
+3 -1
View File
@@ -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<WritePolicyResult> {
// 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
+11 -3
View File
@@ -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<LocalRelay>,
/// Write policy used for promotion-time recovery hooks.
write_policy: Option<Nip34WritePolicy>,
/// Naughty list tracker for git remote domains with persistent errors
git_naughty_list: Arc<NaughtyListTracker>,
}
@@ -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<Purgatory>,
@@ -242,6 +247,7 @@ impl RealSyncContext {
git_data_path: PathBuf,
our_domain: Option<String>,
local_relay: Option<LocalRelay>,
write_policy: Option<Nip34WritePolicy>,
git_naughty_list: Arc<NaughtyListTracker>,
) -> 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<String>,
) -> Result<ProcessResult> {
// 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?;
+287 -31
View File
@@ -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<String> {
}
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;