mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(relay): keep archival snapshots monotonic
NIP-IA snapshots used whole-second timestamps, so rapid archive state
changes could lose to NIP-16's random same-second event-id tie-break and
strand stale state. Advance each snapshot past the current head and retry
from canonical archive state when a concurrent replacement wins or changes
the set before dispatch.
Add a focused archive-to-unarchive regression that verifies the final stored
snapshot advances and carries the canonical empty set.
Co-authored-by: Mongo <5c25403eab7271f9f94ddd4f2b270e8cac2c92e2c830c51877cca6ec974ffb3f@buzz.block.builderlab.xyz>
Signed-off-by: Wes <wesbillman@users.noreply.github.com>
(cherry picked from commit 1ad28a44b9)
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -3124,29 +3124,70 @@ pub async fn publish_nipia_archival_list(
|
||||
tenant: &TenantContext,
|
||||
state: &Arc<AppState>,
|
||||
) -> 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<Tag> = 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<Tag> = 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
|
||||
|
||||
Reference in New Issue
Block a user