From fa50e9f153edca31825ebfd337352c0e2a92a68b Mon Sep 17 00:00:00 2001 From: Ren Amamiya Date: Sun, 27 Sep 2026 02:50:11 +0000 Subject: [PATCH] feat: add event fetching strategy (#24) Reviewed-on: https://git.reya.info/reya/signed/pulls/24 --- AGENTS.md | 2 + CHANGELOG.md | 5 + crates/settings/src/settings.rs | 22 ++ crates/signed_core/src/annotations.rs | 7 - crates/signed_core/src/filters.rs | 40 +- crates/signed_core/src/inbox.rs | 7 +- crates/signed_core/src/lib.rs | 2 - crates/signed_state/src/backend.rs | 194 ++++++---- crates/signed_state/src/checkouts.rs | 50 +-- crates/signed_state/src/inbox.rs | 9 +- crates/signed_state/src/profile.rs | 38 +- crates/signed_state/src/refresh.rs | 4 - crates/signed_state/src/repo.rs | 343 +++++++++++------- crates/signed_state/src/repos.rs | 36 +- crates/workspace/src/views/commit_diff/mod.rs | 4 +- crates/workspace/src/views/inbox.rs | 31 +- .../src/views/pull_requests/detail.rs | 6 +- .../workspace/src/views/pull_requests/new.rs | 15 +- .../src/views/sidebar/grasp_servers.rs | 10 +- crates/workspace/src/views/sidebar/mod.rs | 12 +- .../src/views/sidebar/settings_dialog.rs | 46 ++- 21 files changed, 484 insertions(+), 399 deletions(-) delete mode 100644 crates/signed_core/src/annotations.rs diff --git a/AGENTS.md b/AGENTS.md index 922e8c2..c28ad0a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -81,6 +81,8 @@ Both `cx.spawn` and `cx.background_spawn` return a `Task`, which is a future A task which doesn't do anything but provide a value can be created with `Task::ready(value)`. +Prefer keeping a task in a field over `.detach()` when the work belongs to a view or store. A detached task outlives the view that started it (for example, a repository panel the user has closed), while a task stored in a field is cancelled when that struct drops. A finished `Task` held in a container is not reaped by GPUI - it stays alive until its handle is dropped - so if the container can accumulate many runs, either drop the finished handles before pushing a new one, or hold a single `Option>` and replace it. + ## Elements The `Render` trait is used to render some state into an element tree that is laid out using flexbox layout. An `Entity` where `T` implements `Render` is sometimes called a "view". diff --git a/CHANGELOG.md b/CHANGELOG.md index 0e8d330..abd775e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,20 +8,25 @@ - Show an avatar in each panel's tab, using the repository owner's profile picture when set and a pixel avatar otherwise - Add a Tab Bar setting to hide the previous/next tab buttons, hidden by default +- Add an Event Fetching Strategy setting, fetching a repository's activity from its announced relays only (Curated) or from every maintainer's relays as well (Uncensored, the default) ### Changed - 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 ### Fixed - Render every avatar at one consistent size, where a surrounding border had shrunk pictures by two pixels and the pixel avatar ignored an explicit size - Date repositories from their repository state event, so the explore list and open repository views show the latest push instead of the announcement date +- Fix background tasks outliving their view, so a closed repository panel or pull request view stops fetching and publishing on its own ### Removed +- Remove cover note support, the kind-1624 GitWorkshop and `ngit` extension outside the NIP-34 + ### Deprecated ## v0.1.0-alpha - 2026/09/14 diff --git a/crates/settings/src/settings.rs b/crates/settings/src/settings.rs index a9a7264..bdef5ec 100644 --- a/crates/settings/src/settings.rs +++ b/crates/settings/src/settings.rs @@ -18,6 +18,16 @@ pub enum AppearanceMode { Dark, } +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum EventFetchingStrategy { + /// Only the relays declared in the repository announcement. + Curated, + /// Repository relays plus every maintainer's NIP-65 relays. + #[default] + Uncensored, +} + /// Fields mirror the gpui-component `Theme` surface customized at startup. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(default)] @@ -141,6 +151,7 @@ pub struct CreateRepositorySettings { #[serde(default)] pub struct Settings { pub appearance: AppearanceMode, + pub event_fetching: EventFetchingStrategy, pub theme: ThemeSettings, pub tab_bar: TabBarSettings, pub grasp_servers: GraspServersSettings, @@ -157,6 +168,7 @@ mod tests { fn json_roundtrip_preserves_everything() { let settings = Settings { appearance: AppearanceMode::Dark, + event_fetching: EventFetchingStrategy::Curated, theme: ThemeSettings { radius: 8.0, ..Default::default() @@ -178,9 +190,19 @@ mod tests { serde_json::from_str(r#"{"appearance": "dark", "theme": {"radius": 4.0}}"#).unwrap(); assert_eq!(settings.appearance, AppearanceMode::Dark); assert_eq!(settings.theme.radius, 4.0); + // Unset fields fall back to their defaults, including the fetching strategy. + assert_eq!(settings.event_fetching, EventFetchingStrategy::Uncensored); // The rest of the theme and the other groups keep their defaults. assert_eq!(settings.theme.light_theme, "Signed Light"); assert_eq!(settings.grasp_servers, GraspServersSettings::default()); assert_eq!(settings.create_repository.default_folder, None); } + + #[test] + fn event_fetching_uses_snake_case() { + let json = serde_json::to_string(&EventFetchingStrategy::Curated).unwrap(); + assert_eq!(json, r#""curated""#); + let parsed: EventFetchingStrategy = serde_json::from_str(r#""uncensored""#).unwrap(); + assert_eq!(parsed, EventFetchingStrategy::Uncensored); + } } diff --git a/crates/signed_core/src/annotations.rs b/crates/signed_core/src/annotations.rs deleted file mode 100644 index 0a49935..0000000 --- a/crates/signed_core/src/annotations.rs +++ /dev/null @@ -1,7 +0,0 @@ -use nostr::prelude::*; - -/// GitWorkshop and `ngit` cover-note extension, kind 1624. -/// -/// A markdown note attached to an issue, patch or PR by its author or a maintainer, -/// not part of the NIP-34 draft, read support for interop. -pub const COVER_NOTE_KIND: Kind = Kind::Custom(1624); diff --git a/crates/signed_core/src/filters.rs b/crates/signed_core/src/filters.rs index 678cf50..c8180a0 100644 --- a/crates/signed_core/src/filters.rs +++ b/crates/signed_core/src/filters.rs @@ -2,7 +2,7 @@ use std::time::Duration; use nostr::prelude::*; -use crate::{COVER_NOTE_KIND, RepoAddr}; +use crate::RepoAddr; /// Kinds that make up the activity of a repository. pub const ACTIVITY_KINDS: [Kind; 9] = [ @@ -18,19 +18,18 @@ pub const ACTIVITY_KINDS: [Kind; 9] = [ ]; /// Kinds that notify a user when they tag them via their `p` tag. -pub const NOTIFICATION_KINDS: [Kind; 9] = [ +pub const NOTIFICATION_KINDS: [Kind; 8] = [ Kind::GitIssue, Kind::GitPullRequest, Kind::GitPatch, Kind::GitPullRequestUpdate, - COVER_NOTE_KIND, Kind::GitStatusOpen, Kind::GitStatusApplied, Kind::GitStatusClosed, Kind::GitStatusDraft, ]; -/// Git root kinds that make a comment or cover note count as git activity. +/// Git root kinds that make a comment count as git activity. const GIT_ROOT_KINDS: [Kind; 4] = [ Kind::GitIssue, Kind::GitPatch, @@ -38,6 +37,15 @@ 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 + || 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 @@ -140,13 +148,7 @@ pub fn notifications(me: PublicKey) -> Vec { /// A comment on an unrelated kind is matched too, so results must be filtered /// through [`is_git_activity`] before display. pub fn authored_activity(me: PublicKey) -> Filter { - Filter::new() - .kinds( - ACTIVITY_KINDS - .into_iter() - .chain(std::iter::once(COVER_NOTE_KIND)), - ) - .author(me) + Filter::new().kinds(ACTIVITY_KINDS).author(me) } /// Whether a kind-1111 comment targets a git root, checked via its `K` tag. @@ -155,12 +157,6 @@ fn is_git_comment(event: &Event) -> bool { && tag_kind(event, "K").is_some_and(|kind| GIT_ROOT_KINDS.contains(&kind)) } -/// Whether a kind-1624 cover note targets a git root, checked via its `k` tag. -fn is_git_cover_note(event: &Event) -> bool { - event.kind == COVER_NOTE_KIND - && tag_kind(event, "k").is_some_and(|kind| GIT_ROOT_KINDS.contains(&kind)) -} - /// Whether a status event references a git root, checked via its `k` tag. fn is_git_status(event: &Event) -> bool { tag_kind(event, "k").is_some_and(|kind| GIT_ROOT_KINDS.contains(&kind)) @@ -175,7 +171,7 @@ pub fn is_git_activity(event: &Event) -> bool { | Kind::GitStatusApplied | Kind::GitStatusClosed | Kind::GitStatusDraft => is_git_status(event), - kind => kind == COVER_NOTE_KIND && is_git_cover_note(event), + _ => false, } } @@ -267,17 +263,12 @@ mod tests { } #[test] - fn status_and_cover_note_activity_depend_on_the_lowercase_k_tag() { + fn status_activity_depends_on_the_lowercase_k_tag() { let status = signed( &keys(1), Kind::GitStatusClosed, vec![kind_tag("k", Kind::GitPullRequest)], ); - let cover = signed( - &keys(1), - COVER_NOTE_KIND, - vec![kind_tag("k", Kind::GitPatch)], - ); let unrelated = signed( &keys(1), Kind::GitStatusClosed, @@ -285,7 +276,6 @@ mod tests { ); assert!(is_git_activity(&status)); - assert!(is_git_activity(&cover)); assert!(!is_git_activity(&unrelated)); assert!(!is_git_activity(&signed( &keys(1), diff --git a/crates/signed_core/src/inbox.rs b/crates/signed_core/src/inbox.rs index 582e9b0..c848a29 100644 --- a/crates/signed_core/src/inbox.rs +++ b/crates/signed_core/src/inbox.rs @@ -4,7 +4,7 @@ use std::time::Duration; use nostr::prelude::*; use serde::{Deserialize, Serialize}; -use crate::{COVER_NOTE_KIND, RepoAddr, activity_subject}; +use crate::{RepoAddr, activity_subject}; /// Window before `now` that an advanced cutoff retreats to. const ADVANCE_WINDOW: Duration = Duration::from_secs(3 * 24 * 60 * 60); @@ -121,14 +121,11 @@ impl InboxItem { /// - patch (1617): its `e` parent patch, else itself /// - NIP-22 comment (1111): uppercase `E` root pointer /// - PR update (1619): uppercase `E` -/// - statuses (1630-1633) / cover note (1624): NIP-10 root `e` +/// - statuses (1630-1633): NIP-10 root `e` pub fn notification_root(event: &Event, lookup: &L) -> Option where L: Fn(EventId) -> Option, { - if event.kind == COVER_NOTE_KIND { - return nip10_root_id(event).map(|root| resolve_thread_root(root, lookup)); - } match event.kind { Kind::GitIssue | Kind::GitPullRequest => Some(event.id), Kind::GitPatch => Some(match first_e_id(event) { diff --git a/crates/signed_core/src/lib.rs b/crates/signed_core/src/lib.rs index 2bc3ca6..7232086 100644 --- a/crates/signed_core/src/lib.rs +++ b/crates/signed_core/src/lib.rs @@ -1,5 +1,4 @@ pub mod addr; -pub mod annotations; pub mod deletions; pub mod filters; pub mod inbox; @@ -8,7 +7,6 @@ pub mod state; pub mod status; pub use addr::{RepoAddr, identifier_from_name, repo_addr}; -pub use annotations::COVER_NOTE_KIND; pub use deletions::Deletions; pub use filters::{ NOTIFICATION_KINDS, authored_activity, is_git_activity, notification_comments, notifications, diff --git a/crates/signed_state/src/backend.rs b/crates/signed_state/src/backend.rs index 1c461b4..ccdf0cc 100644 --- a/crates/signed_state/src/backend.rs +++ b/crates/signed_state/src/backend.rs @@ -23,12 +23,7 @@ pub const USER_KEYRING: &str = "Signed Safe Storage"; pub const NOSTR_CONNECT_TIMEOUT: u64 = 60; /// Relays connected at startup, before any user-specific relay config is known. -pub const BOOTSTRAP_RELAYS: [&str; 4] = [ - "wss://relay.primal.net", - "wss://relay.ditto.pub", - "wss://index.ngit.dev", - "wss://profiles.nostr1.com", -]; +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"]; @@ -42,20 +37,15 @@ pub enum BackendEvent { PassphraseRequired, /// The signer changed on login, logout or account switch. SignerChanged, - /// New events were received from a relay and stored in the database. - /// - /// Batched: [`Backend`]'s notification pump coalesces everything a - /// relay delivers within one debounce window into a single event, - /// instead of emitting per-event and making every subscriber debounce - /// the same burst independently. - 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, current: u64, }, - /// An event built locally was signed, broadcast and stored. - Published(Box), Error(String), } @@ -100,14 +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 notifications.next().await { - Some(ClientNotification::Event { event, .. }) => { - pending.push(Update::from_event(&event)); + match next_update(&mut notifications, &mut seen).await { + Some(UpdateEvent::Profile(author)) => { + pending_profiles.insert(author); } - Some(_) => continue, + Some(UpdateEvent::Repo(update)) => pending_repos.push(update), None => break, } @@ -123,28 +115,33 @@ impl Backend { let timer = cx.background_executor().timer(deadline - now); futures::pin_mut!(timer); - let next = notifications.next(); + let next = next_update(&mut notifications, &mut seen); futures::pin_mut!(next); match futures::future::select(next, timer).await { - futures::future::Either::Left(( - Some(ClientNotification::Event { event, .. }), - _, - )) => { - pending.push(Update::from_event(&event)); + 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((Some(_), _)) => continue, 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}"); } } @@ -446,9 +443,6 @@ impl Backend { let client = self.client.clone(); // Initialize directly at the user's chosen destination. - // No mirror is pre-populated: `GitCache::ensure_clone` lazily clones - // from the grasp server the first time the repo detail view needs it, - // exactly like every other repository. let destination = { let dir_name = signed_git::sanitize_path_component(&name); let dir_name = if dir_name.is_empty() { @@ -513,9 +507,7 @@ impl Backend { let builder = announcement.into_event_builder(); let event = builder.finalize_async(&signer).await?; let output = client.send_event(&event).broadcast().await?; - let event = require_relay_accepted(output, event)?; - this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?; - event + require_relay_accepted(output, event)? }; // The state event is the push authorization. Stage it on each @@ -570,14 +562,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing makes it // visible to the other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| { - this.announce_published(state_event.clone(), cx) - }) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } let announcement = Announcement::from_event(&event) @@ -666,9 +654,7 @@ impl Backend { let builder = announcement.into_event_builder(); let event = builder.finalize_async(&signer).await?; let output = client.send_event(&event).broadcast().await?; - let event = require_relay_accepted(output, event)?; - this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?; - event + require_relay_accepted(output, event)? }; let refs = state.refs.clone(); @@ -725,12 +711,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing makes it visible to the other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| this.announce_published(state_event.clone(), cx)) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } } @@ -907,14 +891,10 @@ impl Backend { // Fan the state out to the relays once a git server holds the objects. // Staging already stored the event locally, publishing notifies // the repository views and other relays and clients. - if let Some(state_event) = &outcome.state_event { - if let Err(e) = client.send_event(state_event).broadcast().await { - log::warn!("failed to broadcast repository state: {e}"); - } - this.update(cx, |this, cx| { - this.announce_published(state_event.clone(), cx) - }) - .ok(); + if let Some(state_event) = &outcome.state_event + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); } Ok(outcome) @@ -1086,7 +1066,7 @@ impl Backend { ) .await?; - for url in user_grasp_list_servers(client.clone(), public_key).await? { + for url in user_grasp_list_servers(&client, public_key).await? { client.add_relay(url).and_connect().await.ok(); } @@ -1205,6 +1185,21 @@ impl Backend { .detach(); } + /// Sync filters through the SDK's NIP-65 gossip targeting. + pub fn sync_auto(&mut self, filters: Vec, cx: &mut Context) { + let client = self.client.clone(); + + cx.spawn(async move |_this, _cx| { + for filter in filters { + if let Err(e) = client.sync(filter).await { + log::warn!("gossip relay fetch failed: {e}"); + } + } + Ok::<(), Error>(()) + }) + .detach(); + } + pub fn subscribe_bootstrap(&mut self, filters: Vec, cx: &mut Context) { let client = self.client.clone(); @@ -1222,7 +1217,8 @@ impl Backend { .detach(); } - pub fn sync_bootstrap(&mut self, filter: Filter, cx: &mut Context) { + /// Sync several bootstrap filters in order, within a single task. + pub fn sync_bootstraps(&mut self, filters: Vec, cx: &mut Context) { let client = self.client.clone(); let (tx, mut rx) = SyncProgress::channel(); @@ -1259,8 +1255,19 @@ impl Backend { .detach(); let sync = cx.background_spawn(async move { - let opts = SyncOptions::default().progress(tx); - sync_bootstrap_only(&client, filter, opts).await + let mut first_error = None; + + for filter in filters { + let opts = SyncOptions::default().progress(tx.clone()); + if let Err(error) = sync_bootstrap_only(&client, filter, opts).await { + first_error.get_or_insert(error); + } + } + + match first_error { + Some(error) => Err(error), + None => Ok(()), + } }); cx.spawn(async move |this, cx| { @@ -1285,14 +1292,6 @@ impl Backend { .detach(); } - /// Emit [`BackendEvent::Published`] for cross-store invalidation. - /// - /// Callers publish with `client.send_event(...)` directly, then call this - /// so stores like `RepoListStore` refresh without re-querying the relays. - pub fn announce_published(&self, event: Event, cx: &mut Context) { - cx.emit(BackendEvent::Published(Box::new(event))); - } - /// Publish a NIP-09 deletion for each of `events`, best-effort. /// /// Each target gets its own deletion event: a relay rejecting or @@ -1315,6 +1314,42 @@ 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 { + loop { + match notifications.next().await { + Some(ClientNotification::Message { message, .. }) => { + 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, + None => return None, + } + } +} + /// Sign and send a single NIP-09 deletion request for `event`. async fn retract_event( client: &Client, @@ -1324,16 +1359,14 @@ async fn retract_event( let builder = EventDeletionRequest::new() .id(event.id) .into_event_builder(); + let deletion = builder.finalize_async(signer).await?; client.send_event(&deletion).broadcast().await?; + Ok(()) } /// The event was accepted by at least one relay, or a descriptive error otherwise. -/// -/// The SDK does not treat "accepted by zero relays" as an error on its own: -/// [`SendEventOutput::success`] may be empty while the call still returns `Ok`. -/// This turns that case into an error the caller can surface. pub(crate) fn require_relay_accepted( output: SendEventOutput, event: Event, @@ -1419,6 +1452,7 @@ pub(crate) async fn sync_bootstrap_only( .with(BOOTSTRAP_RELAYS) .opts(opts) .await?; + Ok(output.value) } @@ -1491,7 +1525,7 @@ fn latest_grasp_list_servers(events: Vec) -> Vec { } pub async fn user_grasp_list_servers( - client: Client, + client: &Client, user: PublicKey, ) -> Result, Error> { let events: Vec = client diff --git a/crates/signed_state/src/checkouts.rs b/crates/signed_state/src/checkouts.rs index 67e52ce..7822c01 100644 --- a/crates/signed_state/src/checkouts.rs +++ b/crates/signed_state/src/checkouts.rs @@ -17,15 +17,9 @@ use crate::repos::RepoListStore; const REFRESH_DEBOUNCE: Duration = Duration::from_millis(300); /// How often the statuses are recomputed against the local refs. -/// -/// A commit lands in a checkout long before the remote reconciliation cadence, -/// so this fast pass surfaces ready-to-push and ready-to-contribute checkouts -/// within a second or two. It reads the tracking refs only, no network. const LOCAL_POLL: Duration = Duration::from_secs(2); - /// How often a full pass refreshes the remotes while any repository panel is open. const STATUS_POLL: Duration = Duration::from_secs(15); - /// Remote refresh interval for the `ready to push` badges of the user's own repositories. const PUSH_POLL: Duration = Duration::from_secs(60); @@ -46,10 +40,6 @@ pub struct CheckoutStatus { /// Commit the branch points at, for tip-based PR dedupe. pub head: String, /// What the branch is compared against. - /// For ready-to-contribute statuses, the announced HEAD branch. - /// The fallbacks are `main`, then the first local branch. - /// For ready-to-push statuses, the remote-tracking ref. - /// Unpushed commits are counted against it. /// /// It is `refs/remotes/origin/`, else `origin/HEAD` for new branches. pub base: String, @@ -67,11 +57,6 @@ struct Remembered { } /// Global store of local-checkout associations and per-checkout statuses. -/// -/// Readers (the sidebar rows, the repository panels) observe this store and -/// derive what they display from their own snapshots, so publishing needs no -/// fine-grained entities: the store notifies when a slice changed and each -/// reader re-derives only what it shows. pub struct CheckoutsStore { /// Checkout paths per announced repository. by_repo: HashMap>, @@ -144,6 +129,7 @@ impl CheckoutsStore { this.statuses.clear(); this.push_statuses.clear(); cx.notify(); + this.refresh(cx); } })); @@ -282,9 +268,6 @@ impl CheckoutsStore { } /// Re-resolve the associations and the requested statuses. - /// - /// Requests arriving while a pass runs fold into a follow-up, requests - /// arriving while the debounce timer is pending are dropped. pub fn refresh(&mut self, cx: &mut Context) { if self.debounce_pending || self.refresh.request() != RefreshRequest::Schedule { return; @@ -300,12 +283,6 @@ impl CheckoutsStore { } /// One full resolve and apply cycle, the debounced entry point. - /// - /// Re-resolves the associations from the settings, the scan and the - /// announcements, then recomputes the requested statuses against freshly - /// fetched remotes. Full passes run on every input change and on the - /// remote reconciliation cadence ([`Self::local_tick`]); they also restart - /// the fast local pass. fn run_refresh(&mut self, cx: &mut Context) { self.debounce_pending = false; self.refresh.begin(); @@ -431,10 +408,6 @@ impl CheckoutsStore { } /// Schedule the fast local status pass, unless one is already pending. - /// - /// Every [`LOCAL_POLL`] the pass recomputes the requested statuses against - /// the local refs, with no network, so a new commit in a checkout surfaces in - /// a second or two instead of at the next remote reconciliation. fn schedule_local_pass(&mut self, cx: &mut Context) { if self.local_pending { return; @@ -452,10 +425,6 @@ impl CheckoutsStore { } /// The fast local status pass. - /// - /// Recomputes the statuses against the local refs; when the remote - /// reconciliation cadence elapsed, it runs a full pass instead so pushes - /// made elsewhere do not linger as `to push`. fn local_tick(&mut self, cx: &mut Context) { // Nothing watched: the chain idles out until a new request restarts it. if self.status_requested.is_empty() && self.push_requested.is_empty() { @@ -490,10 +459,6 @@ impl CheckoutsStore { } /// Recompute the requested statuses against the tracking refs only. - /// - /// The refs were last refreshed by a full pass. Comparing against them is - /// enough to pick up new local commits, and skipping the network keeps - /// this pass cheap enough to run every [`LOCAL_POLL`]. fn run_local_statuses(&mut self, cx: &mut Context) { let associations = self.by_repo.clone(); @@ -635,10 +600,6 @@ fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option Option { if signed_git::worktree_dirty(path) { return None; @@ -649,8 +610,6 @@ fn checkout_push_status(path: &Path, fetch: bool) -> Option { let origin = signed_git::origin_url(path).ok().flatten()?; if fetch { - // Refresh the remote heads first. - // Commits made elsewhere or pushed from another machine must not linger as `to push`. signed_git::fetch_repo_refs(path, &[origin], "+refs/heads/*:refs/remotes/origin/*").ok(); } @@ -678,10 +637,6 @@ fn checkout_push_status(path: &Path, fetch: bool) -> Option { } /// Compute the requested statuses against the checkout paths of `associations`. -/// -/// Shared by the full and the local pass. `fetch` refreshes the checkouts' -/// remote heads first, so the full pass sees remote moves; the fast local -/// pass reads the tracking refs only, which is enough to detect local commits. fn compute_statuses( associations: &HashMap>, requested: &[(RepoAddr, Option)], @@ -737,12 +692,14 @@ pub fn pr_proposes_checkout( if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user { return false; } + let branch_matches = pr .tags .iter() .find(|t| t.kind() == "branch-name") .and_then(|t| t.content()) .is_some_and(|name| name == checkout.branch); + // A renamed branch falls back to the proposed tip commit. let tip_matches = pr .tags @@ -750,6 +707,7 @@ pub fn pr_proposes_checkout( .find(|t| t.kind() == "c") .and_then(|t| t.content()) .is_some_and(|tip| tip == checkout.head); + branch_matches || tip_matches } diff --git a/crates/signed_state/src/inbox.rs b/crates/signed_state/src/inbox.rs index 8d66fec..aed4baa 100644 --- a/crates/signed_state/src/inbox.rs +++ b/crates/signed_state/src/inbox.rs @@ -104,11 +104,16 @@ impl Inbox { /// Sign the state with a random key and store it locally. fn persist(&mut self, cx: &mut Context) { - let Some(me) = Backend::global(cx).read(cx).current_user() else { + let backend = Backend::global(cx); + let (me, client) = { + let backend = backend.read(cx); + (backend.current_user(), backend.client()) + }; + + let Some(me) = me else { return; }; - let client = Backend::global(cx).read(cx).client(); let state = self.state.clone(); let task: Task> = cx.background_spawn(async move { diff --git a/crates/signed_state/src/profile.rs b/crates/signed_state/src/profile.rs index 2e25ac5..b05af89 100644 --- a/crates/signed_state/src/profile.rs +++ b/crates/signed_state/src/profile.rs @@ -92,38 +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| match event { - BackendEvent::NostrUpdate(updates) => { - for update in updates - .iter() - .filter(|update| update.kind == Kind::Metadata) - { - this.apply_author(update.author, cx); + let subscription = cx.subscribe(&backend, |this, _backend, event, cx| { + if let BackendEvent::ProfileUpdates(authors) = event { + for author in authors { + this.apply_author(*author, cx); } } - BackendEvent::Published(event) if event.kind == Kind::Metadata => { - let metadata = Metadata::from_json(&event.content).unwrap_or_default(); - this.profiles - .insert(event.pubkey, Profile::new(event.pubkey, metadata)); - cx.notify(); - } - _ => {} }); - // 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}"); } }); @@ -228,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(); @@ -305,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); @@ -327,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/refresh.rs b/crates/signed_state/src/refresh.rs index e527e0f..695cb7c 100644 --- a/crates/signed_state/src/refresh.rs +++ b/crates/signed_state/src/refresh.rs @@ -18,10 +18,6 @@ impl RefreshGate { self.running } - /// A new refresh request arrived. - /// - /// Folded into a follow-up run while one is in flight, otherwise the - /// caller starts the run itself. pub fn request(&mut self) -> RefreshRequest { if self.running { self.dirty = true; diff --git a/crates/signed_state/src/repo.rs b/crates/signed_state/src/repo.rs index 25cdf78..ffd9623 100644 --- a/crates/signed_state/src/repo.rs +++ b/crates/signed_state/src/repo.rs @@ -4,14 +4,16 @@ use std::path::PathBuf; use anyhow::{Error, bail}; use bitcoin_hashes::sha1::Hash as Sha1Hash; -use gpui::{App, AppContext, AsyncApp, Context, SharedString, Subscription, Task, WeakEntity}; +use gpui::{App, AppContext, Context, SharedString, Subscription, Task}; use nostr::event::IntoEventBuilder; use nostr_sdk::prelude::*; +use settings::{EventFetchingStrategy, SettingsStore}; use signed_core::{ Announcement, Deletions, RepoAddr, RepoStatus, filters, parse_state, pull_request_patch, pull_request_patches, }; use signed_git::Nip34Binding; +use signed_nostr::UniversalSigner; use crate::backend::{ Backend, BackendEvent, grasp_base_url, grasp06_prs_url, pr_clone_urls, require_relay_accepted, @@ -19,7 +21,6 @@ use crate::backend::{ }; use crate::checkouts::CheckoutsStore; use crate::git_store::ensure_repo_mirror; -use crate::refresh::{RefreshGate, RefreshRequest}; use crate::repos::RepoListStore; /// Maximum size of one patch event. @@ -30,7 +31,6 @@ const MAX_PATCH_EVENT_BYTES: usize = 60 * 1024; /// Per-repository store. /// /// Holds the announcement, state, issues, patches, PRs, comments and resolved statuses. -/// Always derived from the local database. pub struct RepoStore { /// NIP-34 address. `None` while the repository is local-only. addr: Option, @@ -84,7 +84,12 @@ pub struct RepoStore { /// /// The per-root fetches cover NIP-22 comments and statuses without an `a` tag. root_fetches: HashSet, - refresh: RefreshGate, + /// Maintainers already synced through gossip in Uncensored mode. + /// + /// Avoids re-running the maintainer Auto sync on every refresh. + synced_maintainers: HashSet, + /// In-flight tasks, cancelled when the store drops. + tasks: Vec>>, /// Backend subscription of an announced repository. `None` while local-only. _subscription: Option, } @@ -132,7 +137,8 @@ impl RepoStore { cloning: false, repo_relays: HashSet::new(), root_fetches: HashSet::new(), - refresh: RefreshGate::default(), + synced_maintainers: HashSet::new(), + tasks: Vec::new(), _subscription: Some(subscription), } } @@ -160,7 +166,8 @@ impl RepoStore { cloning: false, repo_relays: HashSet::new(), root_fetches: HashSet::new(), - refresh: RefreshGate::default(), + synced_maintainers: HashSet::new(), + tasks: Vec::new(), _subscription: None, } } @@ -201,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; @@ -218,20 +224,6 @@ impl RepoStore { deletion || coordinate || authored || comment || status }), - BackendEvent::Published(event) => { - let announcement = event.kind == Kind::GitRepoAnnouncement; - let author = event.pubkey == addr.public_key; - let coordinate = event.tags.coordinates().into_iter().any(|c| c == *addr); - - let state = event.kind == Kind::RepoState - && author - && event.tags.identifier().as_deref() == Some(addr.identifier.as_str()); - - let deletion = - event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish; - - coordinate || (announcement && author) || state || deletion - } _ => false, }; @@ -288,9 +280,71 @@ impl RepoStore { self.repo_relays.extend(new.iter().cloned()); let backend = Backend::global(cx); + let filters = Self::repo_filters(&addr); backend.update(cx, |backend, cx| { - backend.connect_repo_relays(new, Self::repo_filters(&addr), cx); + backend.connect_repo_relays(new, filters, cx); + }); + } + + /// Filters the SDK resolves through NIP-65 gossip in Uncensored mode. + fn maintainer_filters(addr: &RepoAddr, maintainers: &[PublicKey]) -> Vec { + let mut pubkeys = maintainers.to_vec(); + // NIP-34 events tag the announcement author, + // which may not be a maintainer for subordinate forks. + if !pubkeys.contains(&addr.public_key) { + pubkeys.push(addr.public_key); + } + + vec![ + // Announcement and state events, including co-maintainer states. + Filter::new() + .kinds([Kind::GitRepoAnnouncement, Kind::RepoState]) + .authors(pubkeys.clone()) + .identifier(addr.identifier.clone()), + // Activity tagging a maintainer, resolved to their read relays. + Filter::new() + .kinds(filters::ACTIVITY_KINDS) + .coordinate(addr) + .pubkeys(pubkeys.clone()), + // Activity authored by a maintainer, resolved to their write relays. + Filter::new() + .kinds(filters::ACTIVITY_KINDS) + .coordinate(addr) + .authors(pubkeys.clone()), + // Deletions authored by a maintainer. + Filter::new() + .kinds([Kind::EventDeletion, Kind::RequestToVanish]) + .authors(pubkeys), + ] + } + + /// In Uncensored mode, sync the maintainer-shaped filters through the SDK's NIP-65 gossip targeting + fn sync_maintainer_relays(&mut self, maintainers: &[PublicKey], cx: &mut Context) { + let strategy = SettingsStore::try_global(cx) + .map(|store| store.read(cx).settings().event_fetching) + .unwrap_or_default(); + if strategy != EventFetchingStrategy::Uncensored { + return; + } + + let Some(addr) = self.addr.clone() else { + return; + }; + + if !maintainers + .iter() + .any(|public_key| !self.synced_maintainers.contains(public_key)) + { + return; + } + self.synced_maintainers.extend(maintainers.iter().copied()); + + let filters = Self::maintainer_filters(&addr, maintainers); + let backend = Backend::global(cx); + + backend.update(cx, |backend, cx| { + backend.sync_auto(filters, cx); }); } @@ -311,11 +365,6 @@ impl RepoStore { if self.addr.is_none() { return; } - - if self.refresh.request() != RefreshRequest::Schedule { - return; - } - self.run_refresh(cx); } @@ -324,8 +373,6 @@ impl RepoStore { return; }; - self.refresh.begin(); - let backend = Backend::global(cx); let client = backend.read(cx).client(); @@ -444,7 +491,7 @@ impl RepoStore { )) }); - cx.spawn(async move |this, cx| { + let task = cx.spawn(async move |this, cx| { let ( announcement, state, @@ -459,14 +506,13 @@ impl RepoStore { Ok(data) => data, Err(e) => { return this.update(cx, |this, cx| { - this.refresh.abort(); this.last_error = Some(e.to_string()); cx.notify(); }); } }; - let again = this.update(cx, |this, cx| { + this.update(cx, |this, cx| { let keep_hint = announcement.is_none() && !this.loaded; let first_pass = !this.loaded; @@ -488,7 +534,6 @@ impl RepoStore { } // The announcement may list relays for this repository's activity. - // Connect to any we have not fetched from yet. let relays = this .announcement .as_ref() @@ -497,6 +542,15 @@ impl RepoStore { this.connect_announced_relays(&relays, cx); + // Uncensored mode also covers the maintainers' NIP-65 relays. + let maintainers = this + .announcement + .as_ref() + .map(Announcement::effective_maintainers) + .unwrap_or_default(); + + this.sync_maintainer_relays(&maintainers, cx); + if let Some((_, head)) = state { this.head = head; } @@ -542,17 +596,12 @@ impl RepoStore { if changed { cx.notify(); } - - this.refresh.finish() })?; - if again { - this.update(cx, |this, cx| this.refresh(cx))?; - } - Ok(()) - }) - .detach(); + }); + + self.tasks.push(task); } /// Resolve the status of a root event, an issue, patch or PR, per NIP-34. @@ -685,9 +734,12 @@ impl RepoStore { }; let backend = Backend::global(cx); - let signer = backend.read(cx).signer(); + let (signer, client, user) = { + let backend = backend.read(cx); + (backend.signer(), backend.client(), backend.current_user()) + }; - let Some(user) = backend.read(cx).current_user() else { + let Some(user) = user else { self.last_error = Some("Sign in to open a pull request".into()); cx.notify(); return; @@ -722,12 +774,12 @@ impl RepoStore { .collect() }; - cx.spawn(async move |this, cx| { - // The PR references the root patch event. - // Viewers can then find the patch without carrying it inline. + let task: Task> = cx.spawn(async move |this, cx| { + // The PR references the root patch, + // viewers can then find the patch without carrying it inline. let root_patch = match publish_patch_series( - &this, - cx, + &client, + &signer, &addr, owner, euc.as_deref(), @@ -752,11 +804,14 @@ impl RepoStore { // Resolve the servers from the author's latest kind-10317 grasp list. // The settings defaults stand in when no list is published or the query fails. let author_servers = { - let query = this.update(cx, |_this, cx| { - let client = Backend::global(cx).read(cx).client(); - user_grasp_list_servers(client, user) - })?; - match cx.background_spawn(query).await { + let query_client = client.clone(); + let published = cx + .background_spawn( + async move { user_grasp_list_servers(&query_client, user).await }, + ) + .await; + + match published { Ok(published) if !published.is_empty() => published, _ => defaults, } @@ -790,24 +845,17 @@ impl RepoStore { }; let builder = this.update(cx, |this, _cx| { - // NIP-34 PRs carry at least one clone URL. - // The tip commit is downloadable from it. - // The author's `/prs/` URLs come first. - // They are author-controlled and most likely alive. - // The announced mirrors follow. - // The list is fixed before signing. - // The pushed ref name embeds the event id. - // Every candidate URL is listed up front. - // Dead URLs are inert, the linked patch stays the source of truth. let prs_urls: Vec = author_targets .iter() .filter_map(|(url, _)| Url::parse(url).ok()) .collect(); + let base_clone = this .announcement .as_ref() .map(|a| a.clone.clone()) .unwrap_or_default(); + let clone = pr_clone_urls(prs_urls, base_clone); let builder = GitPullRequest { @@ -832,8 +880,6 @@ impl RepoStore { })?; // Sign before publishing. - // The tip is pushed to the grasp servers under `refs/nostr/`. - // Nak's convention, readers fetch that ref for the commit behind the `c` tag. let event = cx .background_spawn({ let signer = signer.clone(); @@ -844,20 +890,22 @@ impl RepoStore { if let Some(path) = push_from.as_ref() { let tip = current_commit.to_string(); let reference = format!("refs/nostr/{}", event.id.to_hex()); + let (pushed, failures) = cx .background_spawn({ let path = path.clone(); let tip = tip.clone(); let reference = reference.clone(); // Author servers first, then the announced base grasp servers. - // The extra targets are best-effort redundancy. let targets: Vec<(String, String)> = author_targets .into_iter() .chain(base_targets) .collect(); + async move { let mut failures = Vec::new(); let mut pushed = 0; + for (url, label) in &targets { match signed_git::push_commit_ref( &path, url, &tip, &reference, @@ -866,6 +914,7 @@ impl RepoStore { Err(e) => failures.push(format!("{label}: {e}")), } } + (pushed, failures) } }) @@ -882,8 +931,6 @@ impl RepoStore { } } - let client = this.update(cx, |_this, cx| Backend::global(cx).read(cx).client())?; - let publish_result: Result = async { let output = client.send_event(&event).broadcast().await?; require_relay_accepted(output, event) @@ -900,11 +947,6 @@ impl RepoStore { } }; - this.update(cx, |_this, cx| { - Backend::global(cx) - .update(cx, |backend, cx| backend.announce_published(pr_event.clone(), cx)) - })?; - // A draft PR carries a kind-1633 status event, NIP-34. // Publish it right after the PR event so viewers never show it open. if draft { @@ -914,8 +956,8 @@ impl RepoStore { } Ok(()) - }) - .detach(); + }); + self.tasks.push(task); } /// Generate the patch between `merge_base` and `compare_ref` in `repo_path`, @@ -979,8 +1021,12 @@ impl RepoStore { self.last_warning = None; let backend = Backend::global(cx); + let (user, client, signer) = { + let backend = backend.read(cx); + (backend.current_user(), backend.client(), backend.signer()) + }; - let Some(user) = backend.read(cx).current_user() else { + let Some(user) = user else { self.last_error = Some("Sign in to update the pull request".into()); cx.notify(); return; @@ -1034,19 +1080,21 @@ impl RepoStore { self.not_announced(cx); return; }; + let owner = addr.public_key; let euc = self.announcement.as_ref().and_then(|a| a.euc.clone()); let root = root.clone(); + let clone: Vec = self .announcement .as_ref() .map(|a| a.clone.clone()) .unwrap_or_default(); - cx.spawn(async move |this, cx| { + let task: Task> = cx.spawn(async move |this, cx| { if let Err(e) = publish_patch_series( - &this, - cx, + &client, + &signer, &addr, owner, euc.as_deref(), @@ -1062,7 +1110,7 @@ impl RepoStore { }); } - let builder = this.update(cx, |_this, _cx| { + let builder = { let builder = GitPullRequestUpdate { repository: addr.clone(), pull_request_event: root.id, @@ -1079,13 +1127,7 @@ impl RepoStore { Some(euc) => builder.tag(Tag::parse(["r", euc]).expect("valid r tag")), None => builder, } - })?; - - let (client, signer) = this.update(cx, |_this, cx| { - let backend = Backend::global(cx); - let backend = backend.read(cx); - (backend.client(), backend.signer()) - })?; + }; let publish_result: Result = async { let event = builder.finalize_async(&signer).await?; @@ -1094,25 +1136,16 @@ impl RepoStore { } .await; - match publish_result { - Ok(event) => { - this.update(cx, |_this, cx| { - Backend::global(cx).update(cx, |backend, cx| { - backend.announce_published(event.clone(), cx) - }) - })?; - } - Err(e) => { - return this.update(cx, |this, cx| { - this.last_error = Some(e.to_string()); - cx.notify(); - }); - } + if let Err(e) = publish_result { + return this.update(cx, |this, cx| { + this.last_error = Some(e.to_string()); + cx.notify(); + }); } Ok(()) - }) - .detach(); + }); + self.tasks.push(task); } /// Set the status of a root event. @@ -1238,7 +1271,7 @@ impl RepoStore { } Ok(()) }); - task.detach(); + self.tasks.push(task); } /// The latest announcement of this repository, @@ -1545,24 +1578,16 @@ impl RepoStore { } .await; - match publish_result { - Ok(event) => { - this.update(cx, |_this, cx| { - Backend::global(cx) - .update(cx, |backend, cx| backend.announce_published(event, cx)) - })?; - } - Err(e) => { - this.update(cx, |this, cx| { - this.last_error = Some(e.to_string()); - cx.notify(); - })?; - } + if let Err(e) = publish_result { + this.update(cx, |this, cx| { + this.last_error = Some(e.to_string()); + cx.notify(); + })?; } Ok(()) }); - task.detach(); + self.tasks.push(task); } } @@ -1632,8 +1657,8 @@ fn patch_current_commit(patch: &str) -> Option<&str> { /// Returns the root event, the one a PR references. #[allow(clippy::too_many_arguments)] async fn publish_patch_series( - this: &WeakEntity, - cx: &mut AsyncApp, + client: &Client, + signer: &UniversalSigner, addr: &RepoAddr, owner: PublicKey, euc: Option<&str>, @@ -1641,12 +1666,6 @@ async fn publish_patch_series( first_marker: &str, reply_to: Option, ) -> Result { - let (client, signer) = this.update(cx, |_this, cx| { - let backend = Backend::global(cx); - let backend = backend.read(cx); - (backend.client(), backend.signer()) - })?; - let mut root: Option = None; let mut previous = reply_to; @@ -1659,6 +1678,7 @@ async fn publish_patch_series( }; let mut tags = vec![Tag::coordinate(addr.clone(), None), Tag::public_key(owner)]; + if ix == 0 { if let Ok(tag) = Tag::parse(["t", first_marker]) { tags.push(tag); @@ -1673,28 +1693,25 @@ async fn publish_patch_series( { tags.push(tag); } + if let Some(euc) = euc && let Ok(tag) = Tag::parse(["r", euc]) { tags.push(tag); } + if let Ok(tag) = Tag::parse(["commit", commit]) { tags.push(tag); } + if let Ok(tag) = Tag::parse(["r", commit]) { tags.push(tag); } let builder = EventBuilder::new(Kind::GitPatch, part.clone()).tags(tags); - - let event = builder.finalize_async(&signer).await?; + let event = builder.finalize_async(signer).await?; let output = client.send_event(&event).broadcast().await?; let event = require_relay_accepted(output, event)?; - this.update(cx, |_this, cx| { - Backend::global(cx).update(cx, |backend, cx| { - backend.announce_published(event.clone(), cx) - }) - })?; if root.is_none() { root = Some(event.clone()); @@ -1733,9 +1750,11 @@ fn comment_builder( #[cfg(test)] mod tests { + use std::collections::HashSet; + use nostr_sdk::prelude::*; - use super::{comment_builder, patch_current_commit}; + use super::{RepoStore, comment_builder, patch_current_commit}; #[test] fn parses_format_patch_header() { @@ -1781,4 +1800,60 @@ mod tests { // Signed's own `references_root` must keep matching the comment. assert!(signed_core::references_root(&event, &root.id)); } + + #[test] + fn maintainer_filters_name_owner_and_maintainers() { + let owner = Keys::generate().public_key(); + let maintainer = Keys::generate().public_key(); + let addr = Coordinate::new(Kind::GitRepoAnnouncement, owner).identifier("my-repo"); + + // The owner is not among the maintainers, as on a subordinate fork. + let filters = RepoStore::maintainer_filters(&addr, &[maintainer]); + assert_eq!(filters.len(), 4); + + let expected = HashSet::from([owner, maintainer]); + + // Gossip only resolves pubkeys from `authors` and the lowercase `#p` tag. + let named = |filter: &Filter| -> HashSet { + let authors = filter.authors.iter().flatten().copied(); + let p_tag = filter + .generic_tags + .get(&SingleLetterTag::LOWERCASE_P) + .into_iter() + .flatten() + .filter_map(|value| PublicKey::from_hex(value).ok()); + authors.chain(p_tag).collect() + }; + + // Announcement and state events, scoped to the repository identifier. + let announcement = &filters[0]; + assert_eq!(named(announcement), expected); + assert!( + announcement + .generic_tags + .contains_key(&SingleLetterTag::LOWERCASE_D) + ); + + // Activity filters, scoped to the repository coordinate. + for filter in &filters[1..3] { + assert_eq!(named(filter), expected); + assert!( + filter + .generic_tags + .contains_key(&SingleLetterTag::LOWERCASE_A) + ); + } + + // Deletions, named by author only. + let deletions = &filters[3]; + assert_eq!( + deletions + .authors + .clone() + .unwrap_or_default() + .into_iter() + .collect::>(), + expected + ); + } } diff --git a/crates/signed_state/src/repos.rs b/crates/signed_state/src/repos.rs index f13ee91..be79c6d 100644 --- a/crates/signed_state/src/repos.rs +++ b/crates/signed_state/src/repos.rs @@ -68,32 +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::Published(event) => { - let announcement = event.kind == Kind::GitRepoAnnouncement; - let state = event.kind == Kind::RepoState; - let deletion = - event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish; - - announcement || state || deletion - } BackendEvent::SignerChanged => { this.state_synced_repos.clear(); true } + BackendEvent::RepoUpdates(_) => true, BackendEvent::Synced => true, _ => false, }; @@ -134,10 +113,15 @@ impl RepoListStore { let backend = Backend::global(cx); backend.update(cx, |backend, cx| { - backend.sync_bootstrap(filters::all_announcements(), cx); - backend.sync_bootstrap(filters::all_states(), cx); - // Deletion requests, NIP-09/62, must be known before any announcement is shown. - backend.sync_bootstrap(filters::deletions(), cx); + backend.sync_bootstraps( + vec![ + filters::all_announcements(), + filters::all_states(), + // Deletion requests, NIP-09/62, must be known before any announcement is shown. + filters::deletions(), + ], + cx, + ); }); } diff --git a/crates/workspace/src/views/commit_diff/mod.rs b/crates/workspace/src/views/commit_diff/mod.rs index b907adb..17799c8 100644 --- a/crates/workspace/src/views/commit_diff/mod.rs +++ b/crates/workspace/src/views/commit_diff/mod.rs @@ -303,6 +303,7 @@ pub struct CommitDiffView { loading: bool, error: Option, pane: Entity, + tasks: Vec>>, } impl CommitDiffView { @@ -336,6 +337,7 @@ impl CommitDiffView { loading: true, error: None, pane, + tasks: Vec::new(), } } @@ -383,7 +385,7 @@ impl CommitDiffView { Ok(()) }); - task.detach(); + self.tasks.push(task); } fn render_header(&self, cx: &mut Context) -> AnyElement { diff --git a/crates/workspace/src/views/inbox.rs b/crates/workspace/src/views/inbox.rs index 38fe63a..8f04a67 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::{InboxItem, InboxReadState, RepoAddr}; use signed_state::{ Backend, BackendEvent, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, query_inbox, }; @@ -157,21 +157,8 @@ 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::Synced | BackendEvent::Published(_) => self.refresh(cx), + BackendEvent::RepoUpdates(_) => self.refresh(cx), + BackendEvent::Synced => self.refresh(cx), _ => {} } } @@ -200,12 +187,16 @@ impl InboxView { self.refresh.begin(); let backend = Backend::global(cx); - let Some(me) = backend.read(cx).current_user() else { + let (me, client) = { + let backend = backend.read(cx); + (backend.current_user(), backend.client()) + }; + + let Some(me) = me else { self.refresh.abort(); return; }; - let client = backend.read(cx).client(); let state = self.state.clone(); let work = cx.background_spawn(async move { query_inbox(&client, me, &state).await }); @@ -570,10 +561,6 @@ fn sub_activity(event: &Event, me: Option, cx: &App) -> AnyElement { } fn activity_phrase(kind: Kind) -> &'static str { - if kind == COVER_NOTE_KIND { - return "added a note"; - } - match kind { Kind::GitIssue => "opened an issue", Kind::GitPullRequest => "opened a PR", diff --git a/crates/workspace/src/views/pull_requests/detail.rs b/crates/workspace/src/views/pull_requests/detail.rs index 89dd0e4..d1385c3 100644 --- a/crates/workspace/src/views/pull_requests/detail.rs +++ b/crates/workspace/src/views/pull_requests/detail.rs @@ -79,8 +79,7 @@ pub struct PullRequestDetailView { pane: Entity, commit_item_sizes: Rc>>, commit_scroll_handle: VirtualListScrollHandle, - /// The dock caches item panels, so without this observer a panel opened - /// before the store loaded would stay on its placeholder. + tasks: Vec>>, _subscription: Subscription, } @@ -124,6 +123,7 @@ impl PullRequestDetailView { pane, commit_item_sizes: Rc::new(Vec::new()), commit_scroll_handle: VirtualListScrollHandle::new(), + tasks: Vec::new(), _subscription: subscription, } } @@ -327,7 +327,7 @@ impl PullRequestDetailView { Ok(()) }); - task.detach(); + self.tasks.push(task); } /// Open the diff of `commit_id` in the bottom dock of the area. diff --git a/crates/workspace/src/views/pull_requests/new.rs b/crates/workspace/src/views/pull_requests/new.rs index a3e3ef8..4d29ba1 100644 --- a/crates/workspace/src/views/pull_requests/new.rs +++ b/crates/workspace/src/views/pull_requests/new.rs @@ -68,6 +68,7 @@ pub struct NewPullRequestView { pane: Entity, scroll_handle: VirtualListScrollHandle, item_sizes: Rc>>, + tasks: Vec>>, _subscriptions: Vec, } @@ -323,6 +324,7 @@ impl NewPullRequestView { pane, scroll_handle: VirtualListScrollHandle::new(), item_sizes: Rc::new(Vec::new()), + tasks: Vec::new(), _subscriptions: subscriptions, } } @@ -390,7 +392,8 @@ impl NewPullRequestView { Ok(()) }); - task.detach(); + + self.tasks.push(task); } /// Branches and the current branch are read off the UI thread, then applied. @@ -419,7 +422,8 @@ impl NewPullRequestView { Ok(()) }); - task.detach(); + + self.tasks.push(task); } fn apply_checkout( @@ -636,7 +640,8 @@ impl NewPullRequestView { Ok(()) }); - task.detach(); + + self.tasks.push(task); } #[allow(clippy::too_many_arguments)] @@ -830,7 +835,7 @@ impl NewPullRequestView { Ok(()) }); - task.detach(); + self.tasks.push(task); } fn submit(&mut self, window: &mut Window, cx: &mut Context) { @@ -911,7 +916,7 @@ impl NewPullRequestView { Ok(()) }); - task.detach(); + self.tasks.push(task); } fn open_commit_diff(&mut self, commit_id: &str, window: &mut Window, cx: &mut Context) { diff --git a/crates/workspace/src/views/sidebar/grasp_servers.rs b/crates/workspace/src/views/sidebar/grasp_servers.rs index 2776b4c..1b5d1a9 100644 --- a/crates/workspace/src/views/sidebar/grasp_servers.rs +++ b/crates/workspace/src/views/sidebar/grasp_servers.rs @@ -215,15 +215,19 @@ pub fn load_user_grasp_servers( cx: &mut App, ) { let backend = Backend::global(cx); - let Some(user) = backend.read(cx).current_user() else { + let (user, client) = { + let backend = backend.read(cx); + (backend.current_user(), backend.client()) + }; + + let Some(user) = user else { state.update(cx, |state, _| state.loading_servers = false); return; }; - let client = backend.read(cx).client(); let handle = window.window_handle(); cx.spawn(async move |cx| { - let result = signed_state::user_grasp_list_servers(client, user).await; + let result = signed_state::user_grasp_list_servers(&client, user).await; let _ = cx.update_window(handle, |_, _window, cx| { state.update(cx, |state, _| { diff --git a/crates/workspace/src/views/sidebar/mod.rs b/crates/workspace/src/views/sidebar/mod.rs index fa4eaf9..fd16d89 100644 --- a/crates/workspace/src/views/sidebar/mod.rs +++ b/crates/workspace/src/views/sidebar/mod.rs @@ -79,16 +79,12 @@ impl SidebarPanel { let signer_changed = matches!(event, BackendEvent::SignerChanged); let signer_required = matches!(event, BackendEvent::SignerRequired); - if !signer_changed && !signer_required { - return; - } - if signer_required { this.banner = pick_banner(); cx.notify(); } - if this.refresh(cx) || signer_required { + if this.refresh(cx) || signer_changed { cx.notify(); } })); @@ -131,16 +127,14 @@ impl SidebarPanel { let user = backend.read(cx).current_user(); let (announcements, local_repos, scanning) = { - let repo_list = RepoListStore::global(cx); - let repo_list = repo_list.read(cx); + let repo_list = RepoListStore::global(cx).read(cx); let announcements = user .as_ref() .map(|user| repo_list.announcements_of(user)) .unwrap_or_default(); - let local = LocalReposStore::global(cx); - let local = local.read(cx); + let local = LocalReposStore::global(cx).read(cx); let local_repos = resolve_local_repos(&local.repos, &repo_list.announcements, &announcements); diff --git a/crates/workspace/src/views/sidebar/settings_dialog.rs b/crates/workspace/src/views/sidebar/settings_dialog.rs index 52a129b..e986a3c 100644 --- a/crates/workspace/src/views/sidebar/settings_dialog.rs +++ b/crates/workspace/src/views/sidebar/settings_dialog.rs @@ -18,7 +18,7 @@ use gpui_component::{ ActiveTheme, IconName, IndexPath, Sizable, Theme, ThemeMode, ThemeRegistry, WindowExt, h_flex, v_flex, }; -use settings::{AppearanceMode, Settings, SettingsStore}; +use settings::{AppearanceMode, EventFetchingStrategy, Settings, SettingsStore}; use signed_ui::{SelectOption, setting_block, setting_row}; use super::{normalize_server, server_host}; @@ -52,6 +52,7 @@ fn theme_options(cx: &App) -> (Vec, Vec) { /// Created once when the dialog opens, so control state survives re-renders. struct SettingsControls { appearance: Entity>>, + event_fetching: Entity>>, light_theme: Entity>>, dark_theme: Entity>>, font_size: Entity, @@ -89,6 +90,23 @@ impl SettingsControls { ) }); + let event_fetching_options = vec![ + SelectOption::new("curated", "Curated"), + SelectOption::new("uncensored", "Uncensored"), + ]; + let event_fetching_value = match settings.event_fetching { + EventFetchingStrategy::Curated => "curated", + EventFetchingStrategy::Uncensored => "uncensored", + }; + let event_fetching = cx.new(|cx| { + SelectState::new( + event_fetching_options.clone(), + selected_index(&event_fetching_options, event_fetching_value), + window, + cx, + ) + }); + let (light_options, dark_options) = theme_options(cx); let light_theme = cx.new(|cx| { SelectState::new( @@ -147,6 +165,19 @@ impl SettingsControls { } })); + subscriptions.push(cx.subscribe(&event_fetching, |_, event, cx| { + if let SelectEvent::Confirm(Some(value)) = event { + let strategy = match value.as_ref() { + "curated" => EventFetchingStrategy::Curated, + _ => EventFetchingStrategy::Uncensored, + }; + let store = SettingsStore::global(cx); + store.update(cx, |store, cx| { + store.edit(|settings| settings.event_fetching = strategy, cx); + }); + } + })); + subscriptions.push(cx.subscribe(&light_theme, |_, event, cx| { if let SelectEvent::Confirm(Some(value)) = event { let store = SettingsStore::global(cx); @@ -252,6 +283,7 @@ impl SettingsControls { Self { appearance, + event_fetching, light_theme, dark_theme, font_size, @@ -294,6 +326,8 @@ fn settings_view(controls: &SettingsControls, cx: &mut App) -> impl IntoElement .child(Separator::horizontal()) .child(grasp_servers_section(&settings, controls, cx)) .child(Separator::horizontal()) + .child(event_fetching_section(controls, cx)) + .child(Separator::horizontal()) .child(repositories_section(&settings, controls, cx)) } @@ -306,6 +340,16 @@ fn appearance_section(controls: &SettingsControls, cx: &App) -> impl IntoElement )) } +/// Which relays repository activity is fetched from. +fn event_fetching_section(controls: &SettingsControls, cx: &App) -> impl IntoElement { + v_flex().w_full().gap_3().child(setting_row( + cx, + "Event Fetching Strategy", + "Curated fetches events from relays in the repository's announcement. Uncensored also fetches from every maintainer's relays.", + Select::new(&controls.event_fetching).w_full(), + )) +} + fn theme_section(settings: &Settings, controls: &SettingsControls, cx: &App) -> impl IntoElement { v_flex() .gap_3()