feat: add event fetching strategy #24

Merged
reya merged 10 commits from optimize into master 2026-09-27 02:50:12 +00:00
6 changed files with 75 additions and 63 deletions
Showing only changes of commit c170c4a800 - Show all commits
+10
View File
@@ -38,6 +38,16 @@ const GIT_ROOT_KINDS: [Kind; 4] = [
Kind::GitRepoAnnouncement, 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`. /// Value of the first tag named `name` on `event`.
fn tag_value<'a>(event: &'a Event, name: &str) -> Option<&'a str> { fn tag_value<'a>(event: &'a Event, name: &str) -> Option<&'a str> {
event event
+48 -16
View File
@@ -37,9 +37,10 @@ pub enum BackendEvent {
PassphraseRequired, PassphraseRequired,
/// The signer changed on login, logout or account switch. /// The signer changed on login, logout or account switch.
SignerChanged, SignerChanged,
/// Events received from a relay, including this client's own events echoed /// Kind-0 metadata arrived for these authors; re-read them from the store.
/// back by a relay after publishing. ProfileUpdates(Vec<PublicKey>),
NostrUpdate(Vec<Update>), /// Repository events arrived: announcements, states, activity and deletions.
RepoUpdates(Vec<Update>),
Synced, Synced,
SyncProgress { SyncProgress {
total: u64, total: u64,
@@ -89,12 +90,16 @@ impl Backend {
let pump: Task<Result<(), Error>> = cx.spawn(async move |this, cx| { let pump: Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
let mut notifications = pump_client.notifications(); let mut notifications = pump_client.notifications();
let mut pending: Vec<Update> = Vec::new(); let mut pending_profiles: HashSet<PublicKey> = HashSet::new();
let mut pending_repos: Vec<Update> = Vec::new();
let mut seen: HashSet<EventId> = HashSet::new(); let mut seen: HashSet<EventId> = HashSet::new();
'outer: loop { 'outer: loop {
match next_update(&mut notifications, &mut seen).await { 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, None => break,
} }
@@ -114,18 +119,29 @@ impl Backend {
futures::pin_mut!(next); futures::pin_mut!(next);
match futures::future::select(next, timer).await { 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::Left((None, _)) => break 'outer,
futures::future::Either::Right(_) => break, futures::future::Either::Right(_) => break,
} }
} }
let batch = std::mem::take(&mut pending); let profiles: Vec<PublicKey> = pending_profiles.drain().collect();
let repos = std::mem::take(&mut pending_repos);
if let Err(e) = if let Err(e) = this.update(cx, |_this, cx| {
this.update(cx, |_this, cx| cx.emit(BackendEvent::NostrUpdate(batch))) if !profiles.is_empty() {
{ cx.emit(BackendEvent::ProfileUpdates(profiles));
log::warn!("failed to emit nostr update: {e}"); }
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. /// Await the next relay event from the notification stream.
async fn next_update( async fn next_update(
notifications: &mut (impl futures::Stream<Item = ClientNotification> + Unpin), notifications: &mut (impl futures::Stream<Item = ClientNotification> + Unpin),
seen: &mut HashSet<EventId>, seen: &mut HashSet<EventId>,
) -> Option<Update> { ) -> Option<UpdateEvent> {
loop { loop {
match notifications.next().await { match notifications.next().await {
Some(ClientNotification::Message { message, .. }) => { Some(ClientNotification::Message { message, .. }) => {
if let RelayMessage::Event { event, .. } = *message let RelayMessage::Event { event, .. } = *message else {
&& seen.insert(event.id) continue;
{ };
return Some(Update::from_event(&event));
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, Some(_) => continue,
+13 -16
View File
@@ -92,31 +92,28 @@ impl ProfileStore {
pub(crate) fn new(cx: &mut Context<Self>) -> Self { pub(crate) fn new(cx: &mut Context<Self>) -> Self {
let backend = Backend::global(cx); let backend = Backend::global(cx);
let client = backend.read(cx).client();
let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
if let BackendEvent::NostrUpdate(updates) = event { if let BackendEvent::ProfileUpdates(authors) = event {
for update in updates for author in authors {
.iter() this.apply_author(*author, cx);
.filter(|update| update.kind == Kind::Metadata)
{
this.apply_author(update.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::<PublicKey>();
let entity = cx.entity().downgrade(); let entity = cx.entity().downgrade();
let entity_clone = entity.clone();
let (sender, receiver) = flume::unbounded::<PublicKey>();
cx.spawn(async move |_this, cx| { cx.spawn(async move |_this, cx| {
Self::handle_requests(entity, &client, &receiver, cx).await Self::handle_requests(entity, &client, &receiver, cx).await
}) })
.detach(); .detach();
let weak = cx.entity().downgrade();
cx.defer(move |cx| { 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}"); 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. /// 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<Self>) { fn apply_seen(&mut self, cx: &mut Context<Self>) {
let authors: Vec<PublicKey> = self.seen.borrow().iter().copied().collect(); let authors: Vec<PublicKey> = self.seen.borrow().iter().copied().collect();
@@ -298,15 +293,20 @@ impl ProfileStore {
// The channel has no async timeout, race the receive against a timer. // The channel has no async timeout, race the receive against a timer.
let deadline = Instant::now() + BATCH_TIMEOUT; let deadline = Instant::now() + BATCH_TIMEOUT;
loop { loop {
let now = Instant::now(); let now = Instant::now();
if now >= deadline { if now >= deadline {
break; break;
} }
let timer = cx.background_executor().timer(deadline - now); let timer = cx.background_executor().timer(deadline - now);
futures::pin_mut!(timer); futures::pin_mut!(timer);
let recv = receiver.recv_async(); let recv = receiver.recv_async();
futures::pin_mut!(recv); futures::pin_mut!(recv);
match futures::future::select(recv, timer).await { match futures::future::select(recv, timer).await {
futures::future::Either::Left((Ok(public_key), _)) => { futures::future::Either::Left((Ok(public_key), _)) => {
batch.insert(public_key); batch.insert(public_key);
@@ -320,9 +320,6 @@ impl ProfileStore {
.kind(Kind::Metadata) .kind(Kind::Metadata)
.authors(batch.drain().collect::<Vec<PublicKey>>()); .authors(batch.drain().collect::<Vec<PublicKey>>());
// 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 { match sync_bootstrap_only(client, filter, SyncOptions::default()).await {
Ok(_) => { Ok(_) => {
this.update(cx, |this, cx| this.apply_seen(cx)).ok(); this.update(cx, |this, cx| this.apply_seen(cx)).ok();
+1 -2
View File
@@ -208,8 +208,7 @@ impl RepoStore {
}; };
let relevant = match event { let relevant = match event {
BackendEvent::NostrUpdate(updates) => updates.iter().any(|update| { BackendEvent::RepoUpdates(updates) => updates.iter().any(|update| {
// Deletions may target any event of this repository.
let deletion = let deletion =
update.kind == Kind::EventDeletion || update.kind == Kind::RequestToVanish; update.kind == Kind::EventDeletion || update.kind == Kind::RequestToVanish;
+1 -14
View File
@@ -68,24 +68,11 @@ impl RepoListStore {
let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
let relevant = match event { 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 => { BackendEvent::SignerChanged => {
this.state_synced_repos.clear(); this.state_synced_repos.clear();
true true
} }
BackendEvent::RepoUpdates(_) => true,
BackendEvent::Synced => true, BackendEvent::Synced => true,
_ => false, _ => false,
}; };
+2 -15
View File
@@ -11,7 +11,7 @@ use gpui::{
use gpui_component::button::{Button, ButtonVariants}; use gpui_component::button::{Button, ButtonVariants};
use gpui_component::{ActiveTheme, Icon, IconName, IconNamed, Sizable, StyledExt, h_flex, v_flex}; use gpui_component::{ActiveTheme, Icon, IconName, IconNamed, Sizable, StyledExt, h_flex, v_flex};
use nostr::prelude::{Event, EventId, Kind, PublicKey, Timestamp}; 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::{ use signed_state::{
Backend, BackendEvent, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, query_inbox, 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<Self>) { fn handle_backend_event(&mut self, event: &BackendEvent, cx: &mut Context<Self>) {
match event { match event {
BackendEvent::NostrUpdate(updates) => { BackendEvent::RepoUpdates(_) => self.refresh(cx),
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::Synced => self.refresh(cx), BackendEvent::Synced => self.refresh(cx),
_ => {} _ => {}
} }