fix performance

This commit is contained in:
2026-09-01 07:40:55 +07:00
parent 8dc45d08c0
commit ab2277ee83
11 changed files with 471 additions and 176 deletions
+91 -14
View File
@@ -1,7 +1,9 @@
use std::collections::HashMap;
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::time::Duration;
use std::time::{Duration, Instant};
use anyhow::{Error, anyhow, bail};
use bitcoin_hashes::sha1::Hash as Sha1Hash;
@@ -32,6 +34,12 @@ pub const BOOTSTRAP_RELAYS: [&str; 4] = [
/// Relays used for indexing user's relay list (NIP-65).
pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"];
/// How long an identical fetch/sync request is suppressed after it started.
/// A second panel for the same repository (or the global and per-author
/// list stores at login) doesn't duplicate a sync that just ran; after the
/// window, re-fetching is allowed again so data stays fresh.
const FETCH_DEDUP_WINDOW: Duration = Duration::from_secs(5 * 60);
#[derive(Debug, Clone)]
pub enum BackendEvent {
/// User has no signer configured.
@@ -84,6 +92,10 @@ pub struct Backend {
/// Whether the stored credential is NIP-49 encrypted and a passphrase
/// is still needed to resume the session.
passphrase_required: bool,
/// Fingerprints of recently started fetches/syncs (relay + filter set),
/// so duplicate requests within [`FETCH_DEDUP_WINDOW`] collapse into
/// one. Entries are pruned lazily on the next request.
recent_fetches: HashMap<u64, Instant>,
tasks: Vec<Task<Result<(), Error>>>,
}
@@ -134,6 +146,7 @@ impl Backend {
connected: false,
sync_progress: None,
passphrase_required: false,
recent_fetches: HashMap::new(),
tasks: vec![pump],
};
@@ -1082,11 +1095,28 @@ impl Backend {
}));
}
/// Whether an identical fetch was started within [`FETCH_DEDUP_WINDOW`]
/// and is still recent enough to suppress a duplicate. Records the
/// fingerprint (after pruning expired entries) when returning `false`.
fn fetch_recently_started(&mut self, fingerprint: u64) -> bool {
self.recent_fetches
.retain(|_, started| started.elapsed() < FETCH_DEDUP_WINDOW);
if self.recent_fetches.contains_key(&fingerprint) {
return true;
}
self.recent_fetches.insert(fingerprint, Instant::now());
false
}
/// Connect to relays announced by a repository (NIP-34 `relays` tag) and
/// fetch its events from them: a one-shot auto-closing subscription for
/// `filters`, plus a negentropy sync so issues, patches and PRs stored
/// only on those relays are not missed.
///
/// Deduplicated: an identical request (same relays and filters) started
/// within [`FETCH_DEDUP_WINDOW`] is skipped, so a second panel for the
/// same repository doesn't re-run the fetch.
///
/// Best-effort: failures are logged, not surfaced. The relays stay in
/// the pool, so later publishes for this repository also reach them.
pub fn connect_repo_relays(
@@ -1095,11 +1125,23 @@ impl Backend {
filters: Vec<Filter>,
cx: &mut Context<Self>,
) {
let relay_strs: Vec<&str> = relays.iter().map(|url| url.as_str()).collect();
let fingerprint = fetch_fingerprint(&relay_strs, &filters);
if self.fetch_recently_started(fingerprint) {
log::debug!("skipping duplicate repo relay fetch");
return;
}
let client = self.client.clone();
self.tasks.push(cx.spawn(async move |_this, _cx| {
self.tasks.push(cx.spawn(async move |this, cx| {
if let Err(e) = connect_repo_relays_only(&client, relays, filters).await {
log::warn!("repo relay fetch failed: {e}");
// Allow an immediate retry after a failure.
this.update(cx, |this, _cx| {
this.recent_fetches.remove(&fingerprint);
})
.ok();
}
Ok(())
}));
@@ -1127,7 +1169,17 @@ impl Backend {
/// reconciles the local database with the relays in both directions.
/// Emits [`BackendEvent::SyncProgress`] while running (throttled to
/// whole-percent changes) and [`BackendEvent::Synced`] on completion.
///
/// Deduplicated: an identical sync started within
/// [`FETCH_DEDUP_WINDOW`] is skipped. Observers still see the original
/// sync's progress and completion events.
pub fn sync_bootstrap(&mut self, filter: Filter, cx: &mut Context<Self>) {
let fingerprint = fetch_fingerprint(&BOOTSTRAP_RELAYS, std::slice::from_ref(&filter));
if self.fetch_recently_started(fingerprint) {
log::debug!("skipping duplicate bootstrap sync");
return;
}
let client = self.client.clone();
self.sync_progress = Some((0, 0));
@@ -1185,6 +1237,8 @@ impl Backend {
Err(e) => {
this.update(cx, |this, cx| {
this.sync_progress = None;
// Allow an immediate retry after a failure.
this.recent_fetches.remove(&fingerprint);
cx.emit(BackendEvent::error(e.to_string()))
})?;
}
@@ -1303,6 +1357,20 @@ impl Backend {
}
}
/// Fingerprint of a relay + filter set, for fetch dedup. Relays and
/// filters are sorted first so the fingerprint is order-independent.
fn fetch_fingerprint(relays: &[&str], filters: &[Filter]) -> u64 {
let mut relays: Vec<&str> = relays.to_vec();
relays.sort_unstable();
let mut filters: Vec<&Filter> = filters.iter().collect();
filters.sort_unstable();
let mut hasher = DefaultHasher::new();
relays.hash(&mut hasher);
filters.hash(&mut hasher);
hasher.finish()
}
/// Add the given relays, connect to them, and fetch the filters: a one-shot
/// subscription (auto-closing after EOSE) plus a negentropy sync per filter
/// as a second pass, so events that race with the subscription or relays
@@ -1317,10 +1385,15 @@ async fn connect_repo_relays_only(
return Ok(());
}
let mut added = false;
for url in &relays {
client.add_relay(url).await?;
added |= client.add_relay(url).await?;
}
// Connecting is only needed when the pool grew; connected relays no-op,
// but the call still iterates every relay in the pool.
if added {
client.connect().await;
}
client.connect().await;
let opts = SubscribeAutoCloseOptions::default()
.exit_policy(ReqExitPolicy::ExitOnEOSE)
@@ -1332,17 +1405,21 @@ async fn connect_repo_relays_only(
.collect();
client.subscribe(target).close_on(opts).await?;
for filter in filters {
let sync_opts = SyncOptions::default().initial_timeout(Duration::from_secs(5));
if let Err(e) = client
.sync(filter)
.with(relays.iter())
.opts(sync_opts)
.await
{
log::warn!("repo relay negentropy sync failed: {e}");
// Sync the filters concurrently: each reconciles against every relay
// either way, and a relay without NEG-XX support otherwise serializes
// its initial timeout behind every other filter.
let sync_opts = SyncOptions::default().initial_timeout(Duration::from_secs(5));
let syncs = filters.into_iter().map(|filter| {
let client = &client;
let relays = &relays;
let sync_opts = sync_opts.clone();
async move {
if let Err(e) = client.sync(filter).with(relays.iter()).opts(sync_opts).await {
log::warn!("repo relay negentropy sync failed: {e}");
}
}
}
});
futures::future::join_all(syncs).await;
Ok(())
}
+121 -47
View File
@@ -1,5 +1,5 @@
use std::borrow::Cow;
use std::collections::HashSet;
use std::collections::{HashMap, HashSet};
use std::time::Duration;
use anyhow::Error;
@@ -33,11 +33,21 @@ pub struct RepoStore {
pub pull_requests: Vec<Event>,
/// Comments on issues / PRs, oldest first.
pub comments: Vec<Event>,
statuses: Vec<Event>,
/// Resolved status per root event (issue / patch / PR), recomputed on
/// every refresh so render paths are HashMap lookups instead of
/// scanning all status events per root.
status_by_root: HashMap<EventId, RepoStatus>,
/// Open issue / root PR counts, computed with [`Self::status_by_root`]
/// on every refresh.
open_issue_count: usize,
open_pr_count: usize,
/// Kind-1624 cover notes and kind-1985 label events referencing this
/// repository's roots (ngit / GitWorkshop extensions).
cover_notes: Vec<Event>,
labels: Vec<Event>,
/// Incremented on every applied refresh; views key their derived-data
/// caches to it instead of recomputing on every render.
version: u64,
/// Error of the last action initiated from this store, if any.
pub last_error: Option<String>,
/// Relays announced by this repository (NIP-34 `relays` tag) that we
@@ -111,9 +121,12 @@ impl RepoStore {
patches: Vec::new(),
pull_requests: Vec::new(),
comments: Vec::new(),
statuses: Vec::new(),
status_by_root: HashMap::new(),
open_issue_count: 0,
open_pr_count: 0,
cover_notes: Vec::new(),
labels: Vec::new(),
version: 0,
last_error: None,
repo_relays: HashSet::new(),
root_fetches: HashSet::new(),
@@ -141,8 +154,13 @@ impl RepoStore {
/// deletions targeting it.
fn repo_filters(addr: &RepoAddr) -> Vec<Filter> {
let mut filters = vec![
filters::announcement(addr),
filters::state(addr),
// Announcement and state share author and identifier, so they
// combine into one filter: one fewer negentropy reconciliation
// per relay when fetching from the repo's announced relays.
Filter::new()
.kinds([Kind::GitRepoAnnouncement, Kind::RepoState])
.author(addr.public_key)
.identifier(addr.identifier.clone()),
filters::activity(addr),
];
// Deletion requests (NIP-09/62) must be known before any event of
@@ -293,7 +311,7 @@ impl RepoStore {
.chain(&pull_requests)
.map(|e| e.id);
for root in roots {
for event in db.query(filters::statuses_for(root)).await? {
for event in db.query(filters::statuses_for([root])).await? {
if seen_statuses.insert(event.id) {
statuses.push(event);
}
@@ -312,7 +330,7 @@ impl RepoStore {
.chain(&pull_requests)
.map(|e| e.id);
for root in roots {
for event in db.query(filters::annotations_for(root)).await? {
for event in db.query(filters::annotations_for([root])).await? {
if deletions.is_deleted(&event) {
continue;
}
@@ -331,13 +349,36 @@ impl RepoStore {
sort_newest_first(&mut cover_notes);
sort_newest_first(&mut labels);
// Resolve every root's status once here; render paths do
// HashMap lookups instead of scanning all status events per
// root (quadratic, with an allocation per pair).
let maintainers = announcement
.as_ref()
.map(Announcement::effective_maintainers)
.unwrap_or_default();
let status_by_root =
resolve_statuses(&issues, &patches, &pull_requests, &statuses, &maintainers);
let open_issue_count = issues
.iter()
.filter(|issue| status_of(&status_by_root, issue) == RepoStatus::Open)
.count();
let open_pr_count = pull_requests
.iter()
.filter(|pr| {
pr.kind == Kind::GitPullRequest
&& status_of(&status_by_root, pr) == RepoStatus::Open
})
.count();
Ok::<_, Error>((
announcement,
state,
issues,
patches,
pull_requests,
statuses,
status_by_root,
open_issue_count,
open_pr_count,
comments,
cover_notes,
labels,
@@ -353,7 +394,9 @@ impl RepoStore {
issues,
patches,
pull_requests,
statuses,
status_by_root,
open_issue_count,
open_pr_count,
comments,
cover_notes,
labels,
@@ -389,9 +432,12 @@ impl RepoStore {
this.patches = patches;
this.pull_requests = pull_requests;
this.comments = comments;
this.statuses = statuses;
this.status_by_root = status_by_root;
this.open_issue_count = open_issue_count;
this.open_pr_count = open_pr_count;
this.cover_notes = cover_notes;
this.labels = labels;
this.version = this.version.wrapping_add(1);
// Comments, statuses without an `a` tag, cover notes and
// labels are not addressed to the repository, so fetch them
@@ -411,25 +457,18 @@ impl RepoStore {
.collect();
if !new_roots.is_empty() {
this.root_fetches.extend(new_roots.iter().copied());
let comment_filters = filters::comments_for(new_roots.clone());
let status_filters: Vec<Filter> = new_roots
.iter()
.copied()
.map(filters::statuses_for)
.collect();
let annotation_filters: Vec<Filter> = new_roots
.into_iter()
.map(filters::annotations_for)
.collect();
// Batch the per-root filters: one statuses filter and one
// annotations filter covering all new roots, instead of
// one filter per root (each filter is a separate
// negentropy reconciliation per relay).
let mut root_filters = filters::comments_for(new_roots.clone());
root_filters.push(filters::statuses_for(new_roots.iter().copied()));
root_filters.push(filters::annotations_for(new_roots));
let announced: Vec<RelayUrl> = this.repo_relays.iter().cloned().collect();
let backend = Backend::global(cx);
backend.update(cx, |backend, cx| {
backend.subscribe_bootstrap(comment_filters.clone(), cx);
backend.connect_repo_relays(announced.clone(), comment_filters, cx);
backend.subscribe_bootstrap(status_filters.clone(), cx);
backend.connect_repo_relays(announced.clone(), status_filters, cx);
backend.subscribe_bootstrap(annotation_filters.clone(), cx);
backend.connect_repo_relays(announced, annotation_filters, cx);
backend.subscribe_bootstrap(root_filters.clone(), cx);
backend.connect_repo_relays(announced, root_filters, cx);
});
}
@@ -454,20 +493,17 @@ impl RepoStore {
}));
}
/// Resolve the status of a root event (issue / patch / PR) per NIP-34.
/// Resolve the status of a root event (issue / patch / PR) per NIP-34:
/// a lookup into the map built on the last refresh.
pub fn status_of(&self, root: &Event) -> RepoStatus {
let maintainers = self
.announcement
.as_ref()
.map(Announcement::effective_maintainers)
.unwrap_or_default();
status_of(&self.status_by_root, root)
}
let events = self
.statuses
.iter()
.filter(|e| signed_core::references_root(e, &root.id));
signed_core::resolve_status(events, &root.pubkey, &maintainers)
/// Refresh generation, incremented on every applied refresh. Views use
/// it to key their derived-data caches (filtered lists, counts) so
/// renders that change nothing stay O(1).
pub fn version(&self) -> u64 {
self.version
}
/// The effective cover note of `root` (kind 1624), if any: the latest
@@ -509,21 +545,16 @@ impl RepoStore {
/// Number of open issues: issues whose resolved status is
/// [`RepoStatus::Open`] (issues without status events default to open).
/// Cached on the last refresh.
pub fn issue_count(&self) -> usize {
self.issues
.iter()
.filter(|issue| self.status_of(issue) == RepoStatus::Open)
.count()
self.open_issue_count
}
/// Number of open pull requests: root PR events (not PR updates, whose
/// status is carried by the root) with a resolved status of
/// [`RepoStatus::Open`].
/// [`RepoStatus::Open`]. Cached on the last refresh.
pub fn pull_request_count(&self) -> usize {
self.pull_requests
.iter()
.filter(|pr| pr.kind == Kind::GitPullRequest && self.status_of(pr) == RepoStatus::Open)
.count()
self.open_pr_count
}
/// Whether `user` is the author (owner) of this repository: the public
@@ -856,6 +887,49 @@ where
events.into_iter().max_by_key(|e| e.created_at)
}
/// Status of `root` from the precomputed map; roots without status events
/// default to [`RepoStatus::Open`], like [`signed_core::resolve_status`].
fn status_of(status_by_root: &HashMap<EventId, RepoStatus>, root: &Event) -> RepoStatus {
status_by_root
.get(&root.id)
.copied()
.unwrap_or(RepoStatus::Open)
}
/// Resolve the status of every root event in one pass: status events are
/// indexed by the root they reference (`e`/`E` tag), then each root
/// resolves against its own slice. O(roots + statuses) instead of the
/// O(roots × statuses) of resolving per root on demand.
fn resolve_statuses(
issues: &[Event],
patches: &[Event],
pull_requests: &[Event],
statuses: &[Event],
maintainers: &[PublicKey],
) -> HashMap<EventId, RepoStatus> {
let mut by_root: HashMap<EventId, Vec<&Event>> = HashMap::new();
for event in statuses {
for tag in event.tags.iter() {
if matches!(tag.kind(), "e" | "E")
&& let Some(id) = tag.content().and_then(|hex| EventId::from_hex(hex).ok())
{
by_root.entry(id).or_default().push(event);
}
}
}
issues
.iter()
.chain(patches)
.chain(pull_requests)
.map(|root| {
let events = by_root.get(&root.id).map(Vec::as_slice).unwrap_or(&[]);
let status = signed_core::resolve_status(events.iter().copied(), &root.pubkey, maintainers);
(root.id, status)
})
.collect()
}
fn sort_newest_first(events: &mut [Event]) {
events.sort_by_key(|e| std::cmp::Reverse(e.created_at));
}