diff --git a/crates/signed_state/src/backend.rs b/crates/signed_state/src/backend.rs index 0d6998b..f28088c 100644 --- a/crates/signed_state/src/backend.rs +++ b/crates/signed_state/src/backend.rs @@ -37,15 +37,14 @@ pub enum BackendEvent { PassphraseRequired, /// The signer changed on login, logout or account switch. SignerChanged, - /// New events were received from a relay and stored in the database. + /// Events received from a relay, including this client's own events echoed + /// back by a relay after publishing. NostrUpdate(Vec), Synced, SyncProgress { total: u64, current: u64, }, - /// An event built locally was signed, broadcast and stored. - Published(Box), Error(String), } @@ -91,13 +90,11 @@ impl Backend { let pump: Task> = cx.spawn(async move |this, cx| { let mut notifications = pump_client.notifications(); let mut pending: Vec = Vec::new(); + let mut seen: HashSet = HashSet::new(); 'outer: loop { - match notifications.next().await { - Some(ClientNotification::Event { event, .. }) => { - pending.push(Update::from_event(&event)); - } - Some(_) => continue, + match next_update(&mut notifications, &mut seen).await { + Some(update) => pending.push(update), None => break, } @@ -113,17 +110,11 @@ impl Backend { let timer = cx.background_executor().timer(deadline - now); futures::pin_mut!(timer); - let next = notifications.next(); + let next = next_update(&mut notifications, &mut seen); futures::pin_mut!(next); match futures::future::select(next, timer).await { - futures::future::Either::Left(( - Some(ClientNotification::Event { event, .. }), - _, - )) => { - pending.push(Update::from_event(&event)); - } - futures::future::Either::Left((Some(_), _)) => continue, + futures::future::Either::Left((Some(update), _)) => pending.push(update), futures::future::Either::Left((None, _)) => break 'outer, futures::future::Either::Right(_) => break, } @@ -500,9 +491,7 @@ impl Backend { let builder = announcement.into_event_builder(); let event = builder.finalize_async(&signer).await?; let output = client.send_event(&event).broadcast().await?; - let event = require_relay_accepted(output, event)?; - this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?; - event + require_relay_accepted(output, event)? }; // The state event is the push authorization. Stage it on each @@ -557,14 +546,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing makes it // visible to the other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| { - this.announce_published(state_event.clone(), cx) - }) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } let announcement = Announcement::from_event(&event) @@ -653,9 +638,7 @@ impl Backend { let builder = announcement.into_event_builder(); let event = builder.finalize_async(&signer).await?; let output = client.send_event(&event).broadcast().await?; - let event = require_relay_accepted(output, event)?; - this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?; - event + require_relay_accepted(output, event)? }; let refs = state.refs.clone(); @@ -712,12 +695,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing makes it visible to the other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| this.announce_published(state_event.clone(), cx)) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } } @@ -894,14 +875,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing notifies // the repository views and other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| { - this.announce_published(state_event.clone(), cx) - }) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } Ok(outcome) @@ -1299,14 +1276,6 @@ impl Backend { .detach(); } - /// Emit [`BackendEvent::Published`] for cross-store invalidation. - /// - /// Callers publish with `client.send_event(...)` directly, then call this - /// so stores like `RepoListStore` refresh without re-querying the relays. - pub fn announce_published(&self, event: Event, cx: &mut Context) { - cx.emit(BackendEvent::Published(Box::new(event))); - } - /// Publish a NIP-09 deletion for each of `events`, best-effort. /// /// Each target gets its own deletion event: a relay rejecting or @@ -1329,6 +1298,26 @@ impl Backend { } } +/// Await the next relay event from the notification stream. +async fn next_update( + notifications: &mut (impl futures::Stream + Unpin), + seen: &mut HashSet, +) -> Option { + loop { + match notifications.next().await { + Some(ClientNotification::Message { message, .. }) => { + if let RelayMessage::Event { event, .. } = *message + && seen.insert(event.id) + { + return Some(Update::from_event(&event)); + } + } + Some(_) => continue, + None => return None, + } + } +} + /// Sign and send a single NIP-09 deletion request for `event`. async fn retract_event( client: &Client, @@ -1338,16 +1327,14 @@ async fn retract_event( let builder = EventDeletionRequest::new() .id(event.id) .into_event_builder(); + let deletion = builder.finalize_async(signer).await?; client.send_event(&deletion).broadcast().await?; + Ok(()) } /// The event was accepted by at least one relay, or a descriptive error otherwise. -/// -/// The SDK does not treat "accepted by zero relays" as an error on its own: -/// [`SendEventOutput::success`] may be empty while the call still returns `Ok`. -/// This turns that case into an error the caller can surface. pub(crate) fn require_relay_accepted( output: SendEventOutput, event: Event, diff --git a/crates/signed_state/src/profile.rs b/crates/signed_state/src/profile.rs index 2e25ac5..6474287 100644 --- a/crates/signed_state/src/profile.rs +++ b/crates/signed_state/src/profile.rs @@ -93,8 +93,8 @@ impl ProfileStore { pub(crate) fn new(cx: &mut Context) -> Self { let backend = Backend::global(cx); - let subscription = cx.subscribe(&backend, |this, _backend, event, cx| match event { - BackendEvent::NostrUpdate(updates) => { + let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { + if let BackendEvent::NostrUpdate(updates) = event { for update in updates .iter() .filter(|update| update.kind == Kind::Metadata) @@ -102,13 +102,6 @@ impl ProfileStore { this.apply_author(update.author, cx); } } - BackendEvent::Published(event) if event.kind == Kind::Metadata => { - let metadata = Metadata::from_json(&event.content).unwrap_or_default(); - this.profiles - .insert(event.pubkey, Profile::new(event.pubkey, metadata)); - cx.notify(); - } - _ => {} }); // Fetch requests are queued on a channel, batched into one sync per debounce window. diff --git a/crates/signed_state/src/repo.rs b/crates/signed_state/src/repo.rs index c024078..7d76c7e 100644 --- a/crates/signed_state/src/repo.rs +++ b/crates/signed_state/src/repo.rs @@ -225,20 +225,6 @@ impl RepoStore { deletion || coordinate || authored || comment || status }), - BackendEvent::Published(event) => { - let announcement = event.kind == Kind::GitRepoAnnouncement; - let author = event.pubkey == addr.public_key; - let coordinate = event.tags.coordinates().into_iter().any(|c| c == *addr); - - let state = event.kind == Kind::RepoState - && author - && event.tags.identifier().as_deref() == Some(addr.identifier.as_str()); - - let deletion = - event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish; - - coordinate || (announcement && author) || state || deletion - } _ => false, }; @@ -977,11 +963,6 @@ impl RepoStore { } }; - this.update(cx, |_this, cx| { - Backend::global(cx) - .update(cx, |backend, cx| backend.announce_published(pr_event.clone(), cx)) - })?; - // A draft PR carries a kind-1633 status event, NIP-34. // Publish it right after the PR event so viewers never show it open. if draft { @@ -1171,20 +1152,11 @@ impl RepoStore { } .await; - match publish_result { - Ok(event) => { - this.update(cx, |_this, cx| { - Backend::global(cx).update(cx, |backend, cx| { - backend.announce_published(event.clone(), cx) - }) - })?; - } - Err(e) => { - return this.update(cx, |this, cx| { - this.last_error = Some(e.to_string()); - cx.notify(); - }); - } + if let Err(e) = publish_result { + return this.update(cx, |this, cx| { + this.last_error = Some(e.to_string()); + cx.notify(); + }); } Ok(()) @@ -1622,19 +1594,11 @@ impl RepoStore { } .await; - match publish_result { - Ok(event) => { - this.update(cx, |_this, cx| { - Backend::global(cx) - .update(cx, |backend, cx| backend.announce_published(event, cx)) - })?; - } - Err(e) => { - this.update(cx, |this, cx| { - this.last_error = Some(e.to_string()); - cx.notify(); - })?; - } + if let Err(e) = publish_result { + this.update(cx, |this, cx| { + this.last_error = Some(e.to_string()); + cx.notify(); + })?; } Ok(()) @@ -1767,11 +1731,6 @@ async fn publish_patch_series( let event = builder.finalize_async(&signer).await?; let output = client.send_event(&event).broadcast().await?; let event = require_relay_accepted(output, event)?; - this.update(cx, |_this, cx| { - Backend::global(cx).update(cx, |backend, cx| { - backend.announce_published(event.clone(), cx) - }) - })?; if root.is_none() { root = Some(event.clone()); diff --git a/crates/signed_state/src/repos.rs b/crates/signed_state/src/repos.rs index c6309c8..6b122a9 100644 --- a/crates/signed_state/src/repos.rs +++ b/crates/signed_state/src/repos.rs @@ -82,14 +82,6 @@ impl RepoListStore { is_announcement || is_repo_state } }), - BackendEvent::Published(event) => { - let announcement = event.kind == Kind::GitRepoAnnouncement; - let state = event.kind == Kind::RepoState; - let deletion = - event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish; - - announcement || state || deletion - } BackendEvent::SignerChanged => { this.state_synced_repos.clear(); true diff --git a/crates/workspace/src/views/inbox.rs b/crates/workspace/src/views/inbox.rs index 38fe63a..e68aafc 100644 --- a/crates/workspace/src/views/inbox.rs +++ b/crates/workspace/src/views/inbox.rs @@ -171,7 +171,7 @@ impl InboxView { self.refresh(cx); } } - BackendEvent::Synced | BackendEvent::Published(_) => self.refresh(cx), + BackendEvent::Synced => self.refresh(cx), _ => {} } } diff --git a/crates/workspace/src/views/sidebar/mod.rs b/crates/workspace/src/views/sidebar/mod.rs index fa4eaf9..fd16d89 100644 --- a/crates/workspace/src/views/sidebar/mod.rs +++ b/crates/workspace/src/views/sidebar/mod.rs @@ -79,16 +79,12 @@ impl SidebarPanel { let signer_changed = matches!(event, BackendEvent::SignerChanged); let signer_required = matches!(event, BackendEvent::SignerRequired); - if !signer_changed && !signer_required { - return; - } - if signer_required { this.banner = pick_banner(); cx.notify(); } - if this.refresh(cx) || signer_required { + if this.refresh(cx) || signer_changed { cx.notify(); } })); @@ -131,16 +127,14 @@ impl SidebarPanel { let user = backend.read(cx).current_user(); let (announcements, local_repos, scanning) = { - let repo_list = RepoListStore::global(cx); - let repo_list = repo_list.read(cx); + let repo_list = RepoListStore::global(cx).read(cx); let announcements = user .as_ref() .map(|user| repo_list.announcements_of(user)) .unwrap_or_default(); - let local = LocalReposStore::global(cx); - let local = local.read(cx); + let local = LocalReposStore::global(cx).read(cx); let local_repos = resolve_local_repos(&local.repos, &repo_list.announcements, &announcements);