refactor profile cache

This commit is contained in:
2026-09-27 13:37:16 +07:00
parent bdf0e63279
commit 053a16a5e1
2 changed files with 52 additions and 86 deletions
-2
View File
@@ -21,10 +21,8 @@ use crate::repos::RepoListStore;
pub const USER_KEYRING: &str = "Signed Safe Storage"; pub const USER_KEYRING: &str = "Signed Safe Storage";
/// Timeout for NIP-46 signer responses. /// Timeout for NIP-46 signer responses.
pub const NOSTR_CONNECT_TIMEOUT: u64 = 60; pub const NOSTR_CONNECT_TIMEOUT: u64 = 60;
/// Relays connected at startup, before any user-specific relay config is known. /// 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"]; 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. /// 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"]; pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"];
+50 -82
View File
@@ -5,13 +5,19 @@ use std::time::{Duration, Instant};
use anyhow::Error; use anyhow::Error;
use flume::{Receiver, Sender}; use flume::{Receiver, Sender};
use gpui::{ use gpui::{
App, AppContext, AsyncApp, Context, Entity, Global, SharedString, Subscription, Task, App, AppContext, AsyncApp, Context, Entity, Global, SharedString, Subscription, WeakEntity,
WeakEntity,
}; };
use nostr_sdk::prelude::*; use nostr_sdk::prelude::*;
use utils::shorten_pubkey; use utils::shorten_pubkey;
use crate::backend::{Backend, BackendEvent, sync_bootstrap_only}; use crate::backend::{Backend, BackendEvent, subscribe_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. /// A user profile as plain data for the UI, from the kind-0 metadata.
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
@@ -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. /// Global profile cache.
/// ///
/// Profiles are fetched in batches and kept as plain data. /// Profiles are fetched in batches and kept as plain data.
@@ -92,30 +95,26 @@ 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| {
if let BackendEvent::ProfileUpdates(authors) = event {
for author in authors {
this.apply_author(*author, cx);
}
}
});
let entity = cx.entity().downgrade(); let entity = cx.entity().downgrade();
let entity_clone = entity.clone(); let client = backend.read(cx).client();
let (sender, receiver) = flume::unbounded::<PublicKey>(); let (sender, receiver) = flume::unbounded::<PublicKey>();
cx.spawn(async move |_this, cx| { let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
Self::handle_requests(entity, &client, &receiver, cx).await 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(); .detach();
cx.defer(move |cx| { cx.defer(move |cx| {
if let Err(error) = entity_clone.update(cx, |this, cx| this.load(cx)) { entity.update(cx, |this, cx| this.load(cx)).ok();
log::warn!("profile store dropped before initial load could run: {error}");
}
}); });
Self { Self {
@@ -150,11 +149,9 @@ impl ProfileStore {
let client = backend.read(cx).client(); let client = backend.read(cx).client();
let work = cx.background_spawn(async move { 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?; let events = client.database().query(filter).await?;
// Parse off the main thread.
// Only plain profiles cross back.
let profiles: Vec<Profile> = events let profiles: Vec<Profile> = events
.into_iter() .into_iter()
.map(|event| { .map(|event| {
@@ -166,7 +163,7 @@ impl ProfileStore {
Ok::<_, Error>(profiles) Ok::<_, Error>(profiles)
}); });
let task: Task<Result<(), Error>> = cx.spawn(async move |this, cx| { cx.spawn(async move |this, cx| {
let profiles = work.await?; let profiles = work.await?;
this.update(cx, |this, cx| { this.update(cx, |this, cx| {
@@ -176,51 +173,13 @@ impl ProfileStore {
cx.notify(); cx.notify();
})?; })?;
Ok(()) Ok::<_, Error>(())
}); })
task.detach(); .detach();
} }
fn apply_author(&mut self, public_key: PublicKey, cx: &mut Context<Self>) { /// Re-read the latest metadata of `authors` from the local database in one query.
let backend = Backend::global(cx); fn apply_authors(&mut self, authors: Vec<PublicKey>, cx: &mut Context<Self>) {
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<Result<(), Error>> = 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<Self>) {
let authors: Vec<PublicKey> = self.seen.borrow().iter().copied().collect();
if authors.is_empty() { if authors.is_empty() {
return; return;
} }
@@ -231,9 +190,9 @@ impl ProfileStore {
let work = cx.background_spawn(async move { let work = cx.background_spawn(async move {
let filter = Filter::new().kind(Kind::Metadata).authors(authors); let filter = Filter::new().kind(Kind::Metadata).authors(authors);
let events = client.database().query(filter).await?; let events = client.database().query(filter).await?;
// Pick the latest metadata per author off the main thread. // Pick the latest metadata per author off the main thread.
let mut latest: HashMap<PublicKey, (Timestamp, Metadata)> = HashMap::new(); let mut latest: HashMap<PublicKey, (Timestamp, Metadata)> = HashMap::new();
for event in events { for event in events {
match latest.get(&event.pubkey) { match latest.get(&event.pubkey) {
Some((ts, _)) if *ts >= event.created_at => {} Some((ts, _)) if *ts >= event.created_at => {}
@@ -257,7 +216,7 @@ impl ProfileStore {
Ok::<_, Error>(profiles) Ok::<_, Error>(profiles)
}); });
let task: Task<Result<(), Error>> = cx.spawn(async move |this, cx| { cx.spawn(async move |this, cx| {
let profiles = work.await?; let profiles = work.await?;
this.update(cx, |this, cx| { this.update(cx, |this, cx| {
@@ -267,14 +226,18 @@ impl ProfileStore {
cx.notify(); cx.notify();
})?; })?;
Ok(()) Ok::<_, Error>(())
}); })
task.detach(); .detach();
} }
/// Sync metadata for requested authors in batches, debounced to collect requests. /// Re-read the latest metadata of every requested author from the local database.
/// fn apply_seen(&mut self, cx: &mut Context<Self>) {
/// After each batch, the seen profiles are re-read from the database on the main thread. let authors: Vec<PublicKey> = 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( async fn handle_requests(
this: WeakEntity<ProfileStore>, this: WeakEntity<ProfileStore>,
client: &Client, client: &Client,
@@ -316,15 +279,20 @@ impl ProfileStore {
} }
} }
let authors: Vec<PublicKey> = batch.drain().collect();
for chunk in authors.chunks(REQUEST_CHUNK) {
let filter = Filter::new() let filter = Filter::new()
.kind(Kind::Metadata) .kind(Kind::Metadata)
.authors(batch.drain().collect::<Vec<PublicKey>>()); .authors(chunk.iter().copied());
match sync_bootstrap_only(client, filter, SyncOptions::default()).await { if let Err(e) = subscribe_bootstrap_only(client, vec![filter]).await {
Ok(_) => { log::warn!("profile fetch failed: {e}");
this.update(cx, |this, cx| this.apply_seen(cx)).ok();
} }
Err(e) => log::warn!("profile sync failed: {e}"), }
if this.update(cx, |this, cx| this.apply_seen(cx)).is_err() {
return Ok(());
} }
} }
} }