From c170c4a800c128b60ac04c88bf96bcdede8a8477 Mon Sep 17 00:00:00 2001 From: Ren Amamiya Date: Sun, 27 Sep 2026 08:34:13 +0700 Subject: [PATCH] update per event kind --- crates/signed_core/src/filters.rs | 10 +++++ crates/signed_state/src/backend.rs | 64 +++++++++++++++++++++-------- crates/signed_state/src/profile.rs | 29 ++++++------- crates/signed_state/src/repo.rs | 3 +- crates/signed_state/src/repos.rs | 15 +------ crates/workspace/src/views/inbox.rs | 17 +------- 6 files changed, 75 insertions(+), 63 deletions(-) diff --git a/crates/signed_core/src/filters.rs b/crates/signed_core/src/filters.rs index 678cf50..042a00b 100644 --- a/crates/signed_core/src/filters.rs +++ b/crates/signed_core/src/filters.rs @@ -38,6 +38,16 @@ const GIT_ROOT_KINDS: [Kind; 4] = [ Kind::GitRepoAnnouncement, ]; +/// Kinds that carry repository data: announcements, states, activity and deletions. +pub fn is_repo_kind(kind: Kind) -> bool { + kind == Kind::GitRepoAnnouncement + || kind == Kind::RepoState + || kind == Kind::EventDeletion + || kind == Kind::RequestToVanish + || kind == COVER_NOTE_KIND + || ACTIVITY_KINDS.contains(&kind) +} + /// Value of the first tag named `name` on `event`. fn tag_value<'a>(event: &'a Event, name: &str) -> Option<&'a str> { event diff --git a/crates/signed_state/src/backend.rs b/crates/signed_state/src/backend.rs index f28088c..31f6ab2 100644 --- a/crates/signed_state/src/backend.rs +++ b/crates/signed_state/src/backend.rs @@ -37,9 +37,10 @@ pub enum BackendEvent { PassphraseRequired, /// The signer changed on login, logout or account switch. SignerChanged, - /// Events received from a relay, including this client's own events echoed - /// back by a relay after publishing. - NostrUpdate(Vec), + /// Kind-0 metadata arrived for these authors; re-read them from the store. + ProfileUpdates(Vec), + /// Repository events arrived: announcements, states, activity and deletions. + RepoUpdates(Vec), Synced, SyncProgress { total: u64, @@ -89,12 +90,16 @@ 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 pending_profiles: HashSet = HashSet::new(); + let mut pending_repos: Vec = Vec::new(); let mut seen: HashSet = HashSet::new(); 'outer: loop { match next_update(&mut notifications, &mut seen).await { - Some(update) => pending.push(update), + Some(UpdateEvent::Profile(author)) => { + pending_profiles.insert(author); + } + Some(UpdateEvent::Repo(update)) => pending_repos.push(update), None => break, } @@ -114,18 +119,29 @@ impl Backend { futures::pin_mut!(next); match futures::future::select(next, timer).await { - futures::future::Either::Left((Some(update), _)) => pending.push(update), + futures::future::Either::Left((Some(UpdateEvent::Profile(author)), _)) => { + pending_profiles.insert(author); + } + futures::future::Either::Left((Some(UpdateEvent::Repo(update)), _)) => { + pending_repos.push(update); + } futures::future::Either::Left((None, _)) => break 'outer, futures::future::Either::Right(_) => break, } } - let batch = std::mem::take(&mut pending); + let profiles: Vec = pending_profiles.drain().collect(); + let repos = std::mem::take(&mut pending_repos); - if let Err(e) = - this.update(cx, |_this, cx| cx.emit(BackendEvent::NostrUpdate(batch))) - { - log::warn!("failed to emit nostr update: {e}"); + if let Err(e) = this.update(cx, |_this, cx| { + if !profiles.is_empty() { + cx.emit(BackendEvent::ProfileUpdates(profiles)); + } + if !repos.is_empty() { + cx.emit(BackendEvent::RepoUpdates(repos)); + } + }) { + log::warn!("failed to emit backend update: {e}"); } } @@ -1298,18 +1314,34 @@ impl Backend { } } +/// A relay event the backend routes to a store group. +enum UpdateEvent { + Profile(PublicKey), + Repo(Update), +} + /// Await the next relay event from the notification stream. async fn next_update( notifications: &mut (impl futures::Stream + Unpin), seen: &mut HashSet, -) -> Option { +) -> 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)); + let RelayMessage::Event { event, .. } = *message else { + continue; + }; + + let update = match event.kind { + Kind::Metadata => UpdateEvent::Profile(event.pubkey), + kind if filters::is_repo_kind(kind) => { + UpdateEvent::Repo(Update::from_event(&event)) + } + _ => continue, + }; + + if seen.insert(event.id) { + return Some(update); } } Some(_) => continue, diff --git a/crates/signed_state/src/profile.rs b/crates/signed_state/src/profile.rs index 6474287..b05af89 100644 --- a/crates/signed_state/src/profile.rs +++ b/crates/signed_state/src/profile.rs @@ -92,31 +92,28 @@ impl ProfileStore { pub(crate) fn new(cx: &mut Context) -> Self { let backend = Backend::global(cx); + let client = backend.read(cx).client(); 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) - { - this.apply_author(update.author, cx); + if let BackendEvent::ProfileUpdates(authors) = event { + for author in authors { + this.apply_author(*author, cx); } } }); - // Fetch requests are queued on a channel, batched into one sync per debounce window. - let client = backend.read(cx).client(); - let (sender, receiver) = flume::unbounded::(); let entity = cx.entity().downgrade(); + let entity_clone = entity.clone(); + + let (sender, receiver) = flume::unbounded::(); cx.spawn(async move |_this, cx| { Self::handle_requests(entity, &client, &receiver, cx).await }) .detach(); - let weak = cx.entity().downgrade(); cx.defer(move |cx| { - if let Err(error) = weak.update(cx, |this, cx| this.load(cx)) { + if let Err(error) = entity_clone.update(cx, |this, cx| this.load(cx)) { log::warn!("profile store dropped before initial load could run: {error}"); } }); @@ -221,8 +218,6 @@ impl ProfileStore { } /// Re-read the latest metadata of every requested author from the local database. - /// - /// Used after a sync, which produces no NostrUpdate events. fn apply_seen(&mut self, cx: &mut Context) { let authors: Vec = self.seen.borrow().iter().copied().collect(); @@ -298,15 +293,20 @@ impl ProfileStore { // The channel has no async timeout, race the receive against a timer. let deadline = Instant::now() + BATCH_TIMEOUT; + loop { let now = Instant::now(); + if now >= deadline { break; } + let timer = cx.background_executor().timer(deadline - now); futures::pin_mut!(timer); + let recv = receiver.recv_async(); futures::pin_mut!(recv); + match futures::future::select(recv, timer).await { futures::future::Either::Left((Ok(public_key), _)) => { batch.insert(public_key); @@ -320,9 +320,6 @@ impl ProfileStore { .kind(Kind::Metadata) .authors(batch.drain().collect::>()); - // Negentropy-sync with the bootstrap relays. - // Synced events are written to the database directly, no NostrUpdate. - // Re-apply from the database afterwards. match sync_bootstrap_only(client, filter, SyncOptions::default()).await { Ok(_) => { this.update(cx, |this, cx| this.apply_seen(cx)).ok(); diff --git a/crates/signed_state/src/repo.rs b/crates/signed_state/src/repo.rs index 7d76c7e..124208b 100644 --- a/crates/signed_state/src/repo.rs +++ b/crates/signed_state/src/repo.rs @@ -208,8 +208,7 @@ impl RepoStore { }; let relevant = match event { - BackendEvent::NostrUpdate(updates) => updates.iter().any(|update| { - // Deletions may target any event of this repository. + BackendEvent::RepoUpdates(updates) => updates.iter().any(|update| { let deletion = update.kind == Kind::EventDeletion || update.kind == Kind::RequestToVanish; diff --git a/crates/signed_state/src/repos.rs b/crates/signed_state/src/repos.rs index 6b122a9..be79c6d 100644 --- a/crates/signed_state/src/repos.rs +++ b/crates/signed_state/src/repos.rs @@ -68,24 +68,11 @@ impl RepoListStore { let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { let relevant = match event { - BackendEvent::NostrUpdate(updates) => updates.iter().any(|update| { - // Deletions may target anything we list, always refresh. - if update.kind == Kind::EventDeletion || update.kind == Kind::RequestToVanish { - true - } else if filters::ACTIVITY_KINDS.contains(&update.kind) { - // Activity events are addressed to repos via `a` tags. - // Their author is not the repo owner, always refresh. - true - } else { - let is_announcement = update.kind == Kind::GitRepoAnnouncement; - let is_repo_state = update.kind == Kind::RepoState; - is_announcement || is_repo_state - } - }), BackendEvent::SignerChanged => { this.state_synced_repos.clear(); true } + BackendEvent::RepoUpdates(_) => true, BackendEvent::Synced => true, _ => false, }; diff --git a/crates/workspace/src/views/inbox.rs b/crates/workspace/src/views/inbox.rs index e68aafc..7de3cff 100644 --- a/crates/workspace/src/views/inbox.rs +++ b/crates/workspace/src/views/inbox.rs @@ -11,7 +11,7 @@ use gpui::{ use gpui_component::button::{Button, ButtonVariants}; use gpui_component::{ActiveTheme, Icon, IconName, IconNamed, Sizable, StyledExt, h_flex, v_flex}; use nostr::prelude::{Event, EventId, Kind, PublicKey, Timestamp}; -use signed_core::{COVER_NOTE_KIND, InboxItem, InboxReadState, RepoAddr, filters}; +use signed_core::{COVER_NOTE_KIND, InboxItem, InboxReadState, RepoAddr}; use signed_state::{ Backend, BackendEvent, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, query_inbox, }; @@ -157,20 +157,7 @@ impl InboxView { fn handle_backend_event(&mut self, event: &BackendEvent, cx: &mut Context) { match event { - BackendEvent::NostrUpdate(updates) => { - let relevant = updates.iter().any(|update| { - let is_notification = filters::NOTIFICATION_KINDS.contains(&update.kind); - let is_comment = update.kind == Kind::Comment; - let is_event_deletion = update.kind == Kind::EventDeletion; - let is_request_to_vanish = update.kind == Kind::RequestToVanish; - - is_notification || is_comment || is_event_deletion || is_request_to_vanish - }); - - if relevant { - self.refresh(cx); - } - } + BackendEvent::RepoUpdates(_) => self.refresh(cx), BackendEvent::Synced => self.refresh(cx), _ => {} }