Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions crates/buzz-relay/src/handlers/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1501,6 +1501,14 @@ mod tests {
)
}

fn register_nip43_delta_sub(
state: &AppState,
sub_id: &str,
kind: u16,
) -> (Uuid, mpsc::Receiver<Message>) {
register_global_sub(state, sub_id, Filter::new().kind(Kind::Custom(kind)), None)
}

fn presence_event(status: &str) -> nostr::Event {
EventBuilder::new(Kind::Custom(KIND_PRESENCE_UPDATE as u16), status)
.sign_with_keys(&Keys::generate())
Expand All @@ -1517,6 +1525,16 @@ mod tests {
.expect("sign membership notification")
}

fn nip43_delta_event(kind: u16, target: &Keys) -> nostr::Event {
EventBuilder::new(Kind::Custom(kind), "")
.tags([
nostr::Tag::parse(["-"]).expect("protected tag"),
nostr::Tag::parse(["p", &target.public_key().to_hex()]).expect("p tag"),
])
.sign_with_keys(&Keys::generate())
.expect("sign NIP-43 delta")
}

fn event_from_ws_message(msg: Message) -> nostr::Event {
let Message::Text(text) = msg else {
panic!("expected text ws message");
Expand Down Expand Up @@ -1649,6 +1667,33 @@ mod tests {
);
}

#[tokio::test]
async fn global_nip43_delta_pubsub_events_fan_out_by_kind() {
for kind in [8000_u16, 8001_u16] {
let state = test_state().await;
let target = Keys::generate();
let (_conn_id, mut rx) =
register_nip43_delta_sub(&state, &format!("nip43-{kind}"), kind);
let event = nip43_delta_event(kind, &target);
let event_id = event.id;

fan_out_pubsub_event(
&state,
ChannelEvent {
community_id: buzz_core::tenant::CommunityId::from_uuid(Uuid::nil()),
topic: EventTopic::Global,
event,
},
)
.await;

let delivered =
event_from_ws_message(rx.try_recv().expect("NIP-43 delta delivered"));
assert_eq!(delivered.id, event_id);
assert!(rx.try_recv().is_err(), "NIP-43 delta is delivered once");
}
}

async fn redis_url_if_available() -> Option<String> {
let redis_url =
std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
Expand Down
19 changes: 19 additions & 0 deletions crates/buzz-relay/src/handlers/side_effects.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3012,6 +3012,25 @@ async fn publish_nip43_delta(
return Ok(());
}

// Publish through Redis before local fan-out so subscribers connected to
// other relay pods receive the global NIP-43 delta. Mark the event first so
// this pod suppresses its Redis echo after delivering locally.
state.mark_local_event(tenant.community(), &stored.event.id);
if let Err(e) = state
.pubsub
.publish_event(tenant, EventTopic::Global, &stored.event)
.await
{
state
.local_event_ids
.invalidate(&(tenant.community(), stored.event.id.to_bytes()));
warn!(
target = %target_pubkey_hex,
kind,
"NIP-43 {label} Redis publish failed: {e}"
);
}

// Routed through the guarded send path for uniformity; the access gate
// no-ops for this globally-scoped (channel_id = None) NIP-43 event.
crate::handlers::event::fan_out_event_to_local_subscribers(state, tenant.community(), &stored)
Expand Down
Loading