From 272ddf2674278a098208cb8f7b88eeb22617e2d1 Mon Sep 17 00:00:00 2001 From: Ren Amamiya Date: Sun, 27 Sep 2026 07:20:54 +0000 Subject: [PATCH] chore: connect bootstrap relays on demand (#25) Reviewed-on: https://git.reya.info/reya/signed/pulls/25 --- CHANGELOG.md | 4 +- crates/signed_state/src/backend.rs | 79 ++++++++--------- crates/signed_state/src/profile.rs | 135 +++++++++++------------------ 3 files changed, 90 insertions(+), 128 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index abd775e..1628e72 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,7 +15,9 @@ - Migrate the GPUI foundation to the published `gpui-pre` crates and GPUI Kit 0.6, off the zed and gpui-component git pins - Use the pixel avatar as the single fallback for a missing picture, sized and rounded to match the other avatars - Redesign the dock tab bar, using muted grey active tab, added close buttons, double-click to zoom, and removed panel toolbar -- Connect to fewer relays at startup, keeping only ditto and the git indexer as bootstrap relays +- Connect to bootstrap relays on demand instead of at startup +- Fetch profile metadata in batches of 100 authors, applying every requested profile in a single database query +- Prefetch the 500 most recent profiles at startup instead of 200 ### Fixed diff --git a/crates/signed_state/src/backend.rs b/crates/signed_state/src/backend.rs index ccdf0cc..6ef045e 100644 --- a/crates/signed_state/src/backend.rs +++ b/crates/signed_state/src/backend.rs @@ -21,10 +21,8 @@ use crate::repos::RepoListStore; pub const USER_KEYRING: &str = "Signed Safe Storage"; /// Timeout for NIP-46 signer responses. pub const NOSTR_CONNECT_TIMEOUT: u64 = 60; - /// Relays connected at startup, before any user-specific relay config is known. pub const BOOTSTRAP_RELAYS: [&str; 2] = ["wss://relay.ditto.pub", "wss://index.ngit.dev"]; - /// Relays used to index the user's NIP-65 relay list. pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"]; @@ -151,9 +149,10 @@ impl Backend { pump.detach(); cx.defer(move |cx| { - if let Err(error) = weak.update(cx, |this, cx| this.bootstrap(cx)) { - log::warn!("backend dropped before bootstrap could run: {error}"); - } + weak.update(cx, |this, cx| { + this.restore_session(cx); + }) + .ok(); }); Self { @@ -167,42 +166,6 @@ impl Backend { } } - fn bootstrap(&mut self, cx: &mut Context) { - let client = self.client.clone(); - - let task = cx.background_spawn(async move { - for url in BOOTSTRAP_RELAYS { - client.add_relay(url).await?; - } - - for url in INDEXER_RELAYS { - client - .add_relay(url) - .capabilities(RelayCapabilities::DISCOVERY) - .await?; - } - - client.connect().await; - - Ok::<(), Error>(()) - }); - - let notify_task: Task> = cx.spawn(async move |this, cx| { - match task.await { - Ok(()) => { - this.update(cx, |this, cx| { - this.restore_session(cx); - })?; - } - Err(e) => { - this.update(cx, |_this, cx| cx.emit(BackendEvent::error(e.to_string())))?; - } - } - Ok::<(), Error>(()) - }); - notify_task.detach(); - } - /// Restore the saved session from the Keyring. /// /// - Emits [`BackendEvent::SignerRequired`] when no credential is stored. @@ -227,7 +190,9 @@ impl Backend { let result = async { if content.starts_with("nsec1") { let keys = Keys::new(SecretKey::parse(&content)?); - this.update(cx, |this, cx| this.set_signer(keys, cx))?; + this.update(cx, |this, cx| { + this.set_signer(keys, cx); + })?; } else if content.starts_with("bunker://") { let (base, keys) = extract_master_key(&content); let uri = NostrConnectUri::parse(base)?; @@ -238,14 +203,19 @@ impl Backend { None, )?; signer.auth_url_handler(SignedAuthUrlHandler); - this.update(cx, |this, cx| this.set_signer(signer, cx))?; + + this.update(cx, |this, cx| { + this.set_signer(signer, cx); + })?; } else if content.starts_with("ncryptsec1") { this.update(cx, |this, cx| { this.passphrase_required = true; cx.emit(BackendEvent::PassphraseRequired); })?; } else { - this.update(cx, |_this, cx| cx.emit(BackendEvent::SignerRequired))?; + this.update(cx, |_this, cx| { + cx.emit(BackendEvent::SignerRequired); + })?; } Ok::<_, Error>(()) @@ -1424,10 +1394,29 @@ async fn connect_repo_relays( Ok(()) } +/// Add and connect the startup relays. +async fn ensure_bootstrap_relays(client: &Client) -> Result<(), Error> { + for url in BOOTSTRAP_RELAYS { + client.add_relay(url).and_connect().await?; + } + + for url in INDEXER_RELAYS { + client + .add_relay(url) + .capabilities(RelayCapabilities::DISCOVERY) + .and_connect() + .await?; + } + + Ok(()) +} + pub(crate) async fn subscribe_bootstrap_only( client: &Client, filters: Vec, ) -> Result<(), Error> { + ensure_bootstrap_relays(client).await?; + let opts = SubscribeAutoCloseOptions::default() .exit_policy(ReqExitPolicy::ExitOnEOSE) .timeout(Some(Duration::from_secs(10))); @@ -1447,6 +1436,8 @@ pub(crate) async fn sync_bootstrap_only( filter: Filter, opts: SyncOptions, ) -> Result { + ensure_bootstrap_relays(client).await?; + let output = client .sync(filter) .with(BOOTSTRAP_RELAYS) diff --git a/crates/signed_state/src/profile.rs b/crates/signed_state/src/profile.rs index b05af89..02f84b5 100644 --- a/crates/signed_state/src/profile.rs +++ b/crates/signed_state/src/profile.rs @@ -5,14 +5,20 @@ use std::time::{Duration, Instant}; use anyhow::Error; use flume::{Receiver, Sender}; use gpui::{ - App, AppContext, AsyncApp, Context, Entity, Global, SharedString, Subscription, Task, - WeakEntity, + App, AppContext, AsyncApp, Context, Entity, Global, SharedString, Subscription, WeakEntity, }; use nostr_sdk::prelude::*; use utils::shorten_pubkey; use crate::backend::{Backend, BackendEvent, sync_bootstrap_only}; +/// How long to wait for more requests before firing a batched fetch. +const BATCH_TIMEOUT: Duration = Duration::from_millis(500); +/// Max authors per profile request, keeping each filter within relay limits. +const REQUEST_CHUNK: usize = 100; +/// Recent profiles prefetched at startup and read back from the cache. +const WARM_LIMIT: usize = 500; + /// A user profile as plain data for the UI, from the kind-0 metadata. #[derive(Debug, Clone)] pub struct Profile { @@ -62,9 +68,6 @@ impl Profile { } } -/// How long to wait for more requests before firing a batched sync. -const BATCH_TIMEOUT: Duration = Duration::from_millis(500); - /// Global profile cache. /// /// Profiles are fetched in batches and kept as plain data. @@ -92,30 +95,26 @@ 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::ProfileUpdates(authors) = event { - for author in authors { - this.apply_author(*author, cx); - } - } - }); - let entity = cx.entity().downgrade(); - let entity_clone = entity.clone(); + let client = backend.read(cx).client(); let (sender, receiver) = flume::unbounded::(); - cx.spawn(async move |_this, cx| { - Self::handle_requests(entity, &client, &receiver, cx).await + let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { + if let BackendEvent::ProfileUpdates(authors) = event { + this.apply_authors(authors.clone(), cx); + } + }); + + cx.spawn(async move |this, cx| { + if let Err(e) = Self::handle_requests(this, &client, &receiver, cx).await { + log::error!("Failed to handle requests: {e}"); + } }) .detach(); cx.defer(move |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}"); - } + entity.update(cx, |this, cx| this.load(cx)).ok(); }); Self { @@ -150,11 +149,9 @@ impl ProfileStore { let client = backend.read(cx).client(); let work = cx.background_spawn(async move { - let filter = Filter::new().kind(Kind::Metadata).limit(200); + let filter = Filter::new().kind(Kind::Metadata).limit(WARM_LIMIT); let events = client.database().query(filter).await?; - // Parse off the main thread. - // Only plain profiles cross back. let profiles: Vec = events .into_iter() .map(|event| { @@ -166,7 +163,7 @@ impl ProfileStore { Ok::<_, Error>(profiles) }); - let task: Task> = cx.spawn(async move |this, cx| { + cx.spawn(async move |this, cx| { let profiles = work.await?; this.update(cx, |this, cx| { @@ -176,51 +173,13 @@ impl ProfileStore { cx.notify(); })?; - Ok(()) - }); - task.detach(); + Ok::<_, Error>(()) + }) + .detach(); } - fn apply_author(&mut self, public_key: PublicKey, cx: &mut Context) { - let backend = Backend::global(cx); - let client = backend.read(cx).client(); - - let work = cx.background_spawn(async move { - let filter = Filter::new().kind(Kind::Metadata).author(public_key); - let events = client.database().query(filter).await?; - - // Parse off the main thread. - // Only the profile crosses back. - let profile = events - .into_iter() - .max_by_key(|e| e.created_at) - .map(|event| { - let metadata = Metadata::from_json(event.content).unwrap_or_default(); - Profile::new(event.pubkey, metadata) - }); - - Ok::<_, Error>(profile) - }); - - let task: Task> = cx.spawn(async move |this, cx| { - let profile = work.await?; - - this.update(cx, |this, cx| { - if let Some(profile) = profile { - this.profiles.insert(profile.public_key(), profile); - cx.notify(); - } - })?; - - Ok(()) - }); - task.detach(); - } - - /// Re-read the latest metadata of every requested author from the local database. - fn apply_seen(&mut self, cx: &mut Context) { - let authors: Vec = self.seen.borrow().iter().copied().collect(); - + /// Re-read the latest metadata of `authors` from the local database in one query. + fn apply_authors(&mut self, authors: Vec, cx: &mut Context) { if authors.is_empty() { return; } @@ -231,9 +190,9 @@ impl ProfileStore { let work = cx.background_spawn(async move { let filter = Filter::new().kind(Kind::Metadata).authors(authors); let events = client.database().query(filter).await?; - // Pick the latest metadata per author off the main thread. let mut latest: HashMap = HashMap::new(); + for event in events { match latest.get(&event.pubkey) { Some((ts, _)) if *ts >= event.created_at => {} @@ -257,7 +216,7 @@ impl ProfileStore { Ok::<_, Error>(profiles) }); - let task: Task> = cx.spawn(async move |this, cx| { + cx.spawn(async move |this, cx| { let profiles = work.await?; this.update(cx, |this, cx| { @@ -267,14 +226,18 @@ impl ProfileStore { cx.notify(); })?; - Ok(()) - }); - task.detach(); + Ok::<_, Error>(()) + }) + .detach(); } - /// Sync metadata for requested authors in batches, debounced to collect requests. - /// - /// After each batch, the seen profiles are re-read from the database on the main thread. + /// Re-read the latest metadata of every requested author from the local database. + fn apply_seen(&mut self, cx: &mut Context) { + let authors: Vec = self.seen.borrow().iter().copied().collect(); + self.apply_authors(authors, cx); + } + + /// Fetch metadata for requested authors in batches, debounced to collect requests. async fn handle_requests( this: WeakEntity, client: &Client, @@ -316,15 +279,21 @@ impl ProfileStore { } } - let filter = Filter::new() - .kind(Kind::Metadata) - .authors(batch.drain().collect::>()); + let authors: Vec = batch.drain().collect(); - match sync_bootstrap_only(client, filter, SyncOptions::default()).await { - Ok(_) => { - this.update(cx, |this, cx| this.apply_seen(cx)).ok(); + for chunk in authors.chunks(REQUEST_CHUNK) { + let opts = SyncOptions::default(); + let filter = Filter::new() + .kind(Kind::Metadata) + .authors(chunk.iter().copied()); + + if let Err(e) = sync_bootstrap_only(client, filter, opts).await { + log::warn!("profile fetch failed: {e}"); } - Err(e) => log::warn!("profile sync failed: {e}"), + } + + if this.update(cx, |this, cx| this.apply_seen(cx)).is_err() { + return Ok(()); } } }