diff --git a/crates/buzz-relay/src/handlers/identity_archive.rs b/crates/buzz-relay/src/handlers/identity_archive.rs index 9da920483..8bf4fe260 100644 --- a/crates/buzz-relay/src/handlers/identity_archive.rs +++ b/crates/buzz-relay/src/handlers/identity_archive.rs @@ -512,6 +512,93 @@ mod tests { TenantContext::resolved(CommunityId::from_uuid(id), host) } + #[tokio::test] + async fn archival_snapshot_advances_timestamp_for_rapid_state_replacement() { + let Some(pool) = test_pool().await else { + return; + }; + if sqlx::query("SELECT 1 FROM archived_identities LIMIT 1") + .execute(&pool) + .await + .is_err() + { + return; + } + let Some(state) = test_state(pool.clone()).await else { + return; + }; + let tenant = seed_test_community(&pool).await; + let target_hex = Keys::generate().public_key().to_hex(); + let request_id = "a".repeat(64); + + state + .db + .archive( + tenant.community(), + &target_hex, + "self", + &target_hex, + None, + None, + &request_id, + ) + .await + .expect("archive identity"); + publish_nipia_archival_list(&tenant, &state) + .await + .expect("publish archived snapshot"); + let archived_snapshot = state + .db + .query_events(&EventQuery { + kinds: Some(vec![buzz_core::kind::KIND_IA_ARCHIVED_LIST as i32]), + pubkey: Some(state.relay_keypair.public_key().to_bytes().to_vec()), + global_only: true, + limit: Some(1), + ..EventQuery::for_community(tenant.community()) + }) + .await + .expect("query archived snapshot") + .into_iter() + .next() + .expect("archived snapshot exists"); + + state + .db + .unarchive(tenant.community(), &target_hex) + .await + .expect("unarchive identity"); + publish_nipia_archival_list(&tenant, &state) + .await + .expect("publish unarchived snapshot"); + let final_snapshot = state + .db + .query_events(&EventQuery { + kinds: Some(vec![buzz_core::kind::KIND_IA_ARCHIVED_LIST as i32]), + pubkey: Some(state.relay_keypair.public_key().to_bytes().to_vec()), + global_only: true, + limit: Some(1), + ..EventQuery::for_community(tenant.community()) + }) + .await + .expect("query final snapshot") + .into_iter() + .next() + .expect("final snapshot exists"); + + assert!( + final_snapshot.event.created_at > archived_snapshot.event.created_at, + "replacement snapshots must not rely on random same-second event-id ordering" + ); + assert!( + !final_snapshot.event.tags.iter().any(|tag| { + let fields = tag.as_slice(); + fields.first().map(String::as_str) == Some("p") + && fields.get(1).map(String::as_str) == Some(target_hex.as_str()) + }), + "final snapshot must reflect the canonical empty archive set" + ); + } + #[tokio::test] async fn owner_archive_rejects_stale_request_after_live_kind0_owner_flip() { let Some(pool) = test_pool().await else { diff --git a/crates/buzz-relay/src/handlers/side_effects.rs b/crates/buzz-relay/src/handlers/side_effects.rs index 282ea7765..5d55b07d2 100644 --- a/crates/buzz-relay/src/handlers/side_effects.rs +++ b/crates/buzz-relay/src/handlers/side_effects.rs @@ -3124,29 +3124,70 @@ pub async fn publish_nipia_archival_list( tenant: &TenantContext, state: &Arc, ) -> anyhow::Result<()> { - let archived = state.db.list_archived(tenant.community()).await?; - let relay_pubkey_hex = state.relay_keypair.public_key().to_hex(); + const MAX_REPLACEMENT_ATTEMPTS: usize = 8; + let relay_pubkey = state.relay_keypair.public_key(); + let relay_pubkey_hex = relay_pubkey.to_hex(); - let mut tags: Vec = Vec::with_capacity(archived.len() + 1); - tags.push(Tag::parse(["-"]).map_err(|e| anyhow::anyhow!("failed to build '-' tag: {e}"))?); + // A concurrent archive mutation can race between reading the current head and + // replacing it. Rebuild from canonical state on rejection so an older snapshot + // can never strand the final archive set. + for _ in 0..MAX_REPLACEMENT_ATTEMPTS { + let archived = state.db.list_archived(tenant.community()).await?; + let mut tags: Vec = Vec::with_capacity(archived.len() + 1); + tags.push(Tag::parse(["-"]).map_err(|e| anyhow::anyhow!("failed to build '-' tag: {e}"))?); - for identity in &archived { - tags.push( - Tag::parse(["p", &identity.pubkey]) - .map_err(|e| anyhow::anyhow!("failed to build p tag: {e}"))?, - ); - } + for identity in &archived { + tags.push( + Tag::parse(["p", &identity.pubkey]) + .map_err(|e| anyhow::anyhow!("failed to build p tag: {e}"))?, + ); + } - let event = EventBuilder::new(Kind::Custom(KIND_IA_ARCHIVED_LIST as u16), "") - .tags(tags) - .sign_with_keys(&state.relay_keypair) - .map_err(|e| anyhow::anyhow!("failed to sign kind:{KIND_IA_ARCHIVED_LIST}: {e}"))?; + // NIP-16 resolves same-second replacements by event id. Force this + // canonical snapshot strictly past the current head instead of letting a + // rapid archive→unarchive randomly preserve the stale archive state. + let now = nostr::Timestamp::now().as_secs(); + let previous = state + .db + .query_events(&buzz_db::event::EventQuery { + kinds: Some(vec![KIND_IA_ARCHIVED_LIST as i32]), + pubkey: Some(relay_pubkey.to_bytes().to_vec()), + limit: Some(1), + global_only: true, + ..buzz_db::event::EventQuery::for_community(tenant.community()) + }) + .await?; + let created_at = previous + .first() + .map(|event| (event.event.created_at.as_secs() + 1).max(now)) + .unwrap_or(now); + + let event = EventBuilder::new(Kind::Custom(KIND_IA_ARCHIVED_LIST as u16), "") + .tags(tags) + .custom_created_at(nostr::Timestamp::from(created_at)) + .sign_with_keys(&state.relay_keypair) + .map_err(|e| anyhow::anyhow!("failed to sign kind:{KIND_IA_ARCHIVED_LIST}: {e}"))?; + + let (stored, was_inserted) = state + .db + .replace_addressable_event(tenant.community(), &event, None) + .await?; + if !was_inserted { + continue; + } + + let current_archived = state.db.list_archived(tenant.community()).await?; + let snapshot_is_current = + archived + .iter() + .map(|identity| identity.pubkey.as_str()) + .eq(current_archived + .iter() + .map(|identity| identity.pubkey.as_str())); + if !snapshot_is_current { + continue; + } - let (stored, was_inserted) = state - .db - .replace_addressable_event(tenant.community(), &event, None) - .await?; - if was_inserted { dispatch_persistent_event( tenant, state, @@ -3156,13 +3197,16 @@ pub async fn publish_nipia_archival_list( None, ) .await; + info!( + archived_count = archived.len(), + "NIP-IA archived identities list published" + ); + return Ok(()); } - info!( - archived_count = archived.len(), - "NIP-IA archived identities list published" - ); - Ok(()) + anyhow::bail!( + "failed to publish kind:{KIND_IA_ARCHIVED_LIST} after {MAX_REPLACEMENT_ATTEMPTS} concurrent replacements" + ) } /// NIP-DV: publish the relay-signed, per-viewer DM visibility snapshot for