diff --git a/crates/signed_state/src/backend.rs b/crates/signed_state/src/backend.rs index df799bf..140203b 100644 --- a/crates/signed_state/src/backend.rs +++ b/crates/signed_state/src/backend.rs @@ -1,31 +1,28 @@ -use std::collections::{HashMap, HashSet}; +use std::collections::HashSet; use std::path::{Path, PathBuf}; use std::str::FromStr; use std::time::{Duration, Instant}; use anyhow::{Error, anyhow, bail}; use bitcoin_hashes::sha1::Hash as Sha1Hash; -use gpui::{App, AppContext, BackgroundExecutor, Context, Entity, EventEmitter, Global, Task}; +use gpui::{App, AppContext, AsyncApp, Context, Entity, EventEmitter, Global, Task, WeakEntity}; use nostr::event::IntoEventBuilder; use nostr::nips::nip19::Nip19Coordinate; use nostr_connect::prelude::*; -use nostr_sdk::client::SyncSummary; use nostr_sdk::prelude::*; -use signed_core::{Announcement, Filters, RepoAddr, RepoState, filters}; +use signed_core::{Announcement, Filters, RepoAddr, filters}; use signed_git::{GitCache, Repo}; use signed_nostr::{SignedAuthUrlHandler, UniversalSigner, Update}; -use crate::git_store::repo_mirror_path; +use crate::bootstrap::{subscribe_bootstrap_only, sync_bootstrap_only, user_grasp_list_servers}; +use crate::git_store::Mirrors; use crate::inbox::Inbox; +use crate::push::{GraspPush, PushOutcome, grasp_base_url, grasp_clone_url}; 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"]; const PUMP_DEBOUNCE: Duration = Duration::from_millis(200); @@ -89,7 +86,7 @@ impl Backend { let mut seen: HashSet = HashSet::new(); 'outer: loop { - match next_update(&mut notifications, &mut seen).await { + match UpdateEvent::next(&mut notifications, &mut seen).await { Some(UpdateEvent::Profile(author)) => { pending_profiles.insert(author); } @@ -109,7 +106,7 @@ impl Backend { let timer = cx.background_executor().timer(deadline - now); futures::pin_mut!(timer); - let next = next_update(&mut notifications, &mut seen); + let next = UpdateEvent::next(&mut notifications, &mut seen); futures::pin_mut!(next); match futures::future::select(next, timer).await { @@ -338,20 +335,16 @@ impl Backend { .map(|url| RelayUrl::parse(url).expect("valid relay URL")) .collect(); - let client = this.client.clone(); - let signer = this.signer.clone(); + let pusher = GraspPush::new(this.client.clone(), this.signer.clone()); for builder in [ RelayList::new(relays).into_event_builder(), metadata, GitUserGraspList { grasp_servers }.into_event_builder(), ] { - let client = client.clone(); - let signer = signer.clone(); - cx.spawn(async move |_this, _cx| { - publish_best_effort(&client, &signer, builder).await - }) - .detach(); + let pusher = pusher.clone(); + cx.spawn(async move |_this, _cx| pusher.publish_best_effort(builder).await) + .detach(); } })?; @@ -446,94 +439,32 @@ impl Backend { let commit = work.await?; let commit_sha = Sha1Hash::from_str(&commit).map_err(|_| anyhow!("invalid id"))?; - // The nostr client queues events until each relay is connected. - for url in &servers { - client.add_relay(url).and_connect().await.ok(); - } - - // The state event is the push authorization. It must be accepted before the push below. - let announcement = GitRepositoryAnnouncement { - id: repo_id.clone(), - name: Some(name.clone()), - description: (!description.is_empty()).then_some(description.clone()), - web: Vec::new(), - clone: servers - .iter() - .filter_map(|relay| grasp_clone_url(relay, &owner, &repo_id)) - .collect(), - relays: servers.clone(), - euc: Some(commit_sha), - maintainers: Vec::new(), - }; - let signer = this.update(cx, |this, _cx| this.signer.clone())?; - let event = { - let builder = announcement.into_event_builder(); - let event = builder.finalize_async(&signer).await?; - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event)? - }; + let announcement = repository_announcement( + &repo_id, + &name, + &description, + &owner, + &servers, + Some(commit_sha), + ); - // The state event is the push authorization. Stage it on each - // grasp server's relay, then push the initial commit. - // Creation fails only when no server accepted the push, the announcement - // is then retracted so the repository is not left announced without content. - let refs = vec![("refs/heads/main".to_owned(), commit)]; - - let push = cx.background_spawn({ - let client = client.clone(); - let signer = signer.clone(); - let destination = destination.clone(); - let owner = owner.clone(); - let repo_id = repo_id.clone(); - let servers = servers.clone(); - let refs = refs.clone(); - let executor = cx.background_executor().clone(); - async move { - push_staged_to_grasps( - &client, - &signer, - &repo_id, - &refs, - Some("main"), - &destination, - &owner, - &servers, - &executor, - |path, base, owner, repo_id| { - Repo::open(path)?.push_main(base, owner, repo_id) - }, - ) - .await - } - }); - - let outcome = push.await; - - if outcome.accepted() == 0 { - // The announcement is already published. - // Retract it so the repository is not left announced without content. - this.update(cx, |this, cx| { - this.retract_events(std::slice::from_ref(&event), cx); - }) - .ok(); - - return Err(anyhow!( - "The repository was announced, but the push to every grasp server failed: {}. \ - The announcement has been retracted", - outcome.failure_summary() - )); - } - - // 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 - && let Err(e) = client.send_event(state_event).broadcast().await - { - log::warn!("failed to broadcast repository state: {e}"); - } + let event = announce_repository_and_push( + &this, + &client, + &signer, + announcement, + &repo_id, + &owner, + &servers, + vec![("refs/heads/main".to_owned(), commit)], + Some("main".to_owned()), + &destination, + cx, + |path, base, owner, repo_id| Repo::open(path)?.push_main(base, owner, repo_id), + ) + .await?; let announcement = Announcement::from_event(&event) .ok_or_else(|| anyhow!("failed to parse announcement"))?; @@ -596,97 +527,28 @@ impl Backend { }); let (state, euc) = work.await?; - // The nostr client queues events until each relay is connected. - for url in &servers { - client.add_relay(url).and_connect().await.ok(); - } - - // The state event is the push authorization. It must be accepted before the push below. - let announcement = GitRepositoryAnnouncement { - id: repo_id.clone(), - name: Some(name.clone()), - description: (!description.is_empty()).then_some(description.clone()), - web: Vec::new(), - clone: servers - .iter() - .filter_map(|relay| grasp_clone_url(relay, &owner, &repo_id)) - .collect(), - relays: servers.clone(), - euc: euc.and_then(|commit| Sha1Hash::from_str(&commit).ok()), - maintainers: Vec::new(), - }; - let signer = this.update(cx, |this, _cx| this.signer.clone())?; - let event = { - let builder = announcement.into_event_builder(); - let event = builder.finalize_async(&signer).await?; - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event)? - }; + let euc = euc.and_then(|commit| Sha1Hash::from_str(&commit).ok()); + let announcement = + repository_announcement(&repo_id, &name, &description, &owner, &servers, euc); - let refs = state.refs.clone(); - let head = state.head.clone(); - - // The state event is the push authorization. Stage it on each - // grasp server's relay, then push every branch and tag. The push - // fails only when no server accepted it. The announcement is then - // retracted so the repository is not left announced without content. - // An empty repository has no state to stage and nothing to push. - if !refs.is_empty() { - let push = cx.background_spawn({ - let client = client.clone(); - let signer = signer.clone(); - let path = path.clone(); - let owner = owner.clone(); - let repo_id = repo_id.clone(); - let servers = servers.clone(); - let refs = refs.clone(); - let head = head.clone(); - let executor = cx.background_executor().clone(); - async move { - push_staged_to_grasps( - &client, - &signer, - &repo_id, - &refs, - head.as_deref(), - &path, - &owner, - &servers, - &executor, - |path, base, owner, repo_id| { - Repo::open(path)?.push_all(base, owner, repo_id) - }, - ) - .await - } - }); - let outcome = push.await; - - if outcome.accepted() == 0 { - // The announcement is already published. Retract it so - // the repository is not left announced without content. - this.update(cx, |this, cx| { - this.retract_events(std::slice::from_ref(&event), cx); - }) - .ok(); - - return Err(anyhow!( - "The repository was announced, but the push to every grasp server failed: {}. \ - The announcement has been retracted", - outcome.failure_summary() - )); - } - - // 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 - && let Err(e) = client.send_event(state_event).broadcast().await - { - log::warn!("failed to broadcast repository state: {e}"); - } - } + // The state event is the push authorization. It must be accepted before the push below. + let event = announce_repository_and_push( + &this, + &client, + &signer, + announcement, + &repo_id, + &owner, + &servers, + state.refs.clone(), + state.head.clone(), + &path, + cx, + |path, base, owner, repo_id| Repo::open(path)?.push_all(base, owner, repo_id), + ) + .await?; // Point `origin` at the first grasp server so later pushes have a target. if let Some(base) = servers.first().and_then(grasp_base_url) { @@ -734,7 +596,7 @@ impl Backend { announcement: Announcement, cx: &mut Context, ) -> Task> { - let path = repo_mirror_path(&announcement.addr()); + let path = Mirrors::path(&announcement.addr()); self.push_repo_from(announcement, path, None, cx) } @@ -827,8 +689,7 @@ impl Backend { PushOutcome::default() } else { let push = cx.background_spawn({ - let client = client.clone(); - let signer = signer.clone(); + let pusher = GraspPush::new(client.clone(), signer.clone()); let path = path.clone(); let owner = owner.clone(); let repo_id = repo_id.clone(); @@ -837,21 +698,20 @@ impl Backend { let head = head.clone(); let executor = cx.background_executor().clone(); async move { - push_staged_to_grasps( - &client, - &signer, - &repo_id, - &refs, - head.as_deref(), - &path, - &owner, - &relays, - &executor, - |path, base, owner, repo_id| { - Repo::open(path)?.push_all(base, owner, repo_id) - }, - ) - .await + pusher + .push_staged_to_grasps( + &repo_id, + &refs, + head.as_deref(), + &path, + &owner, + &relays, + &executor, + |path, base, owner, repo_id| { + Repo::open(path)?.push_all(base, owner, repo_id) + }, + ) + .await } }); push.await @@ -1035,12 +895,32 @@ impl Backend { filters: Vec, cx: &mut Context, ) { + if relays.is_empty() || filters.is_empty() { + return; + } + let client = self.client.clone(); cx.spawn(async move |_this, _cx| { - if let Err(e) = connect_repo_relays(&client, relays, filters).await { - log::warn!("repo relay fetch failed: {e}"); + let connected: Result<(), Error> = async { + for url in relays.iter() { + client.add_relay(url).and_connect().await?; + } + Ok(()) } + .await; + + if let Err(e) = connected { + log::warn!("repo relay fetch failed: {e}"); + return Ok::<(), Error>(()); + } + + for filter in filters.into_iter() { + if let Err(e) = client.sync(filter).with(relays.iter()).await { + log::warn!("repo relay negentropy sync failed: {e}"); + } + } + Ok::<(), Error>(()) }) .detach(); @@ -1122,15 +1002,13 @@ impl Backend { /// Each target gets its own deletion event: a relay rejecting or /// dropping one does not affect the others. fn retract_events(&mut self, events: &[Event], cx: &mut Context) { - let client = self.client.clone(); - let signer = self.signer.clone(); + let pusher = GraspPush::new(self.client.clone(), self.signer.clone()); for event in events.iter().cloned() { - let client = client.clone(); - let signer = signer.clone(); + let pusher = pusher.clone(); cx.spawn(async move |_this, _cx| { - if let Err(e) = retract_event(&client, &signer, &event).await { + if let Err(e) = pusher.retract_event(&event).await { log::warn!("failed to retract event {}: {e}", event.id); } }) @@ -1139,548 +1017,154 @@ impl Backend { } } +/// The announcement of one repository, as published to the relays. +fn repository_announcement( + repo_id: &str, + name: &str, + description: &str, + owner: &str, + servers: &[RelayUrl], + euc: Option, +) -> GitRepositoryAnnouncement { + GitRepositoryAnnouncement { + id: repo_id.to_owned(), + name: Some(name.to_owned()), + description: (!description.is_empty()).then(|| description.to_owned()), + web: Vec::new(), + clone: servers + .iter() + .filter_map(|relay| grasp_clone_url(relay, owner, repo_id)) + .collect(), + relays: servers.to_vec(), + euc, + maintainers: Vec::new(), + } +} + +/// Announce a repository, stage the announcement on each grasp server's relay +#[allow(clippy::too_many_arguments)] +async fn announce_repository_and_push( + backend: &WeakEntity, + client: &Client, + signer: &UniversalSigner, + announcement: GitRepositoryAnnouncement, + repo_id: &str, + owner: &str, + servers: &[RelayUrl], + refs: Vec<(String, String)>, + head: Option, + path: &Path, + cx: &mut AsyncApp, + push: impl Fn(&Path, &str, &str, &str) -> Result<(), Error> + Send + 'static, +) -> Result { + // The nostr client queues events until each relay is connected. + for url in servers { + client.add_relay(url).and_connect().await.ok(); + } + + // The state event is the push authorization. It must be accepted before the push below. + let event = { + let builder = announcement.into_event_builder(); + let event = builder.finalize_async(signer).await?; + let output = client.send_event(&event).broadcast().await?; + GraspPush::require_relay_accepted(output, event)? + }; + + // The state event is the push authorization. Stage it on each grasp + // server's relay, then push the git data. The push fails only when no + // server accepted it. The announcement is then retracted so the + // repository is not left announced without content. + let outcome = if refs.is_empty() { + PushOutcome::default() + } else { + let pusher = GraspPush::new(client.clone(), signer.clone()); + let repo_id = repo_id.to_owned(); + let owner = owner.to_owned(); + let servers = servers.to_vec(); + let head = head.clone(); + let path = path.to_path_buf(); + let executor = cx.background_executor().clone(); + cx.background_spawn(async move { + pusher + .push_staged_to_grasps( + &repo_id, + &refs, + head.as_deref(), + &path, + &owner, + &servers, + &executor, + push, + ) + .await + }) + .await + }; + + if outcome.accepted() == 0 { + // The announcement is already published. Retract it so the + // repository is not left announced without content. + backend + .update(cx, |backend, cx| { + backend.retract_events(std::slice::from_ref(&event), cx); + }) + .ok(); + + return Err(anyhow!( + "The repository was announced, but the push to every grasp server failed: {}. \ + The announcement has been retracted", + outcome.failure_summary() + )); + } + + // 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 + && let Err(e) = client.send_event(state_event).broadcast().await + { + log::warn!("failed to broadcast repository state: {e}"); + } + + Ok(event) +} + /// 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; - }; +impl UpdateEvent { + /// Await the next relay event from the notification stream. + async fn next( + 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)) + 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); } - _ => 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, - signer: &UniversalSigner, - event: &Event, -) -> Result<(), Error> { - 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. -pub(crate) fn require_relay_accepted( - output: SendEventOutput, - event: Event, -) -> Result { - if output.success.is_empty() && !output.failed.is_empty() { - let reasons = output - .failed - .values() - .cloned() - .collect::>() - .join(", "); - bail!("event not accepted by any relay: {reasons}"); - } - - Ok(event) -} - -/// Sign and broadcast `builder`, logging rather than surfacing failures. -/// -/// Used for best-effort identity bootstrap events, where a relay hiccup -/// should not block sign-up. -async fn publish_best_effort(client: &Client, signer: &UniversalSigner, builder: EventBuilder) { - let result: Result<(), Error> = async { - let event = builder.finalize_async(signer).await?; - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event)?; - Ok(()) - } - .await; - - if let Err(e) = result { - log::warn!("failed to publish identity bootstrap event: {e}"); - } -} - -async fn connect_repo_relays( - client: &Client, - relays: Vec, - filters: Vec, -) -> Result<(), Error> { - if relays.is_empty() || filters.is_empty() { - return Ok(()); - } - - for url in relays.iter() { - client.add_relay(url).and_connect().await?; - } - - for filter in filters.into_iter() { - if let Err(e) = client.sync(filter).with(relays.iter()).await { - log::warn!("repo relay negentropy sync failed: {e}"); - } - } - - 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))); - - let target: HashMap<&str, Vec> = BOOTSTRAP_RELAYS - .iter() - .map(|relay| (*relay, filters.clone())) - .collect(); - - client.subscribe(target).close_on(opts).await?; - - Ok(()) -} - -pub(crate) async fn sync_bootstrap_only( - client: &Client, - filter: Filter, - opts: SyncOptions, -) -> Result { - ensure_bootstrap_relays(client).await?; - - let output = client - .sync(filter) - .with(BOOTSTRAP_RELAYS) - .opts(opts) - .await?; - - Ok(output.value) -} - -/// Base URL of a grasp server, `https://`. -/// -/// `ws://` grasp servers use `http://`, like ngit. -pub(crate) fn grasp_base_url(relay: &RelayUrl) -> Option { - // `domain()` drops the port. - let parsed = Url::parse(relay.as_str()).ok()?; - let host = parsed.host_str()?; - let port = parsed.port().map(|p| format!(":{p}")).unwrap_or_default(); - // `ws://` grasp servers, e.g. local dev relays, speak plain HTTP. - let scheme = if relay.scheme().is_secure() { - "https" - } else { - "http" - }; - Some(format!("{scheme}://{host}{port}")) -} - -fn grasp_clone_url(relay: &RelayUrl, owner: &str, repo_id: &str) -> Option { - let base = grasp_base_url(relay)?; - Url::parse(&format!("{base}/{owner}/{repo_id}.git")).ok() -} - -/// GRASP-06 contributor namespace URL of a pull request tip. -pub(crate) fn grasp06_prs_url(base_url: &str, npub: &str, repo_id: &str) -> String { - format!("{base_url}/prs/{npub}/{repo_id}.git") -} - -/// Assemble the `clone` URLs of a pull request. -/// -/// The author's GRASP-06 `/prs/` URLs come first. -pub(crate) fn pr_clone_urls(prs_urls: Vec, base_clone_urls: Vec) -> Vec { - let mut seen = std::collections::HashSet::new(); - let mut urls = Vec::new(); - for url in prs_urls.into_iter().chain(base_clone_urls) { - if seen.insert(url.to_string()) { - urls.push(url); - } - } - urls -} - -/// The `g` tag servers of one kind-10317 grasp list event, in tag order. -fn grasp_list_servers(event: &Event) -> Vec { - event - .tags - .iter() - .filter(|tag| tag.kind() == "g") - .filter_map(|tag| tag.content()) - .filter_map(|url| RelayUrl::parse(url).ok()) - .collect() -} - -/// Grasp servers of the newest kind-10317 grasp list among `events`. -fn latest_grasp_list_servers(events: Vec) -> Vec { - events - .into_iter() - .max_by_key(|event| event.created_at) - .map(|event| grasp_list_servers(&event)) - .unwrap_or_default() -} - -pub async fn user_grasp_list_servers( - client: &Client, - user: PublicKey, -) -> Result, Error> { - let events: Vec = client - .database() - .query(Filters::grasp_list(user)) - .await? - .into_iter() - .collect(); - Ok(latest_grasp_list_servers(events)) -} - -const GRASP_PUSH_ATTEMPTS: usize = 3; - -/// Pause before re-staging a state event after a transient denial. -const GRASP_RETRY_DELAY: Duration = Duration::from_secs(1); - -#[derive(Debug, Clone)] -pub struct GraspServerResult { - pub relay: RelayUrl, - /// `None` when the server accepted the data, the reason otherwise. - pub reason: Option, -} - -impl GraspServerResult { - fn ok(relay: RelayUrl) -> Self { - Self { - relay, - reason: None, - } - } - - fn failed(relay: RelayUrl, reason: impl Into) -> Self { - Self { - relay, - reason: Some(reason.into()), - } - } -} - -#[derive(Debug, Clone, Default)] -pub struct PushOutcome { - /// Per-server results, in the order the servers were listed. - pub servers: Vec, - /// The newest state event a grasp relay accepted for this push, if any. - /// - /// Broadcast to the other relays once a git server holds the data. - pub state_event: Option, -} - -impl PushOutcome { - pub fn accepted(&self) -> usize { - self.servers - .iter() - .filter(|server| server.reason.is_none()) - .count() - } - - fn failing(&self) -> impl Iterator { - self.servers.iter().filter(|server| server.reason.is_some()) - } - - pub fn failure_summary(&self) -> String { - self.failing() - .map(|server| { - let reason = - utils::flatten_whitespace(server.reason.as_deref().unwrap_or("unknown error")); - format!("{}: {reason}", server.relay) - }) - .collect::>() - .join("; ") - } - - /// A warning for a push only some grasp servers accepted. - /// - /// `None` when every server accepted the push or nothing was pushed. - pub fn partial_warning(&self) -> Option { - let accepted = self.accepted(); - if self.servers.is_empty() || accepted == self.servers.len() { - return None; - } - Some(format!( - "Pushed to {accepted} of {} grasp servers: {}. Republish to sync.", - self.servers.len(), - self.failure_summary() - )) - } -} - -/// Reasons a push attempt should be retried with a freshly staged state -/// event and a fresh git advertisement. -/// -/// Two families are retried: -/// -/// - **Purgatory denials**: the grasp server sends these when the state -/// event for the push has not reached its purgatory yet. Re-staging a -/// fresh event resolves them. -/// - **Stale advertisement races**: `git receive-pack` compares each ref -/// update against the value it advertised when the push started. The grasp -/// server's own background sync can move a ref in between - typically by -/// aligning the repository to a parked state event once the objects of an -/// earlier attempt land - so the compare-and-swap fails with `cannot lock -/// ref` / `incorrect old value provided`. A retry against the fresh -/// advertisement converges, and when the race is lost the pushed data is -/// usually already on the server (see `is_stale_advertisement_race` and -/// the convergence probe in `push_staged_to_grasps`). -/// -/// Other rejections are not retried. -fn is_transient_grasp_denial(stderr: &str) -> bool { - let error = stderr.to_lowercase(); - [ - "no state events in purgatory", - "no matching state event", - "doesn't match push", - "none from authorized publishers", - "no repository announcement found", - "cannot lock ref", - "incorrect old value provided", - ] - .iter() - .any(|marker| error.contains(marker)) -} - -/// A push rejected because `git receive-pack`'s compare-and-swap lost to the -/// grasp server's own background ref alignment: the ref moved between this -/// push's advertisement and its ref transaction (`cannot lock ref ... is at -/// ... but expected ...` / `incorrect old value provided`). The pushed data -/// is usually already on the server by then. -fn is_stale_advertisement_race(stderr: &str) -> bool { - let error = stderr.to_lowercase(); - error.contains("cannot lock ref") || error.contains("incorrect old value provided") -} - -/// Keep `event` as the push's fan-out state event when it is newer than the -/// current one. All staged events carry the same refs; the newest timestamp -/// wins on the relays. -fn keep_newest(state_event: &mut Option, event: Event) { - if state_event - .as_ref() - .is_none_or(|current| event.created_at > current.created_at) - { - *state_event = Some(event); - } -} - -/// Sign a fresh kind `30618` state event for the push. -/// -/// `last_created_at` is the timestamp of the previous event signed for this push. -/// Retries within the same second get the next second: a grasp relay -/// treats a same-id resend as a duplicate and does not re-run its ingest, -/// so an identical resend cannot re-park a state event lost from its purgatory. -async fn sign_state_event( - signer: &UniversalSigner, - repo_id: &str, - refs: &[(String, String)], - head: Option<&str>, - last_created_at: u64, -) -> Result<(Event, u64), String> { - let now = Timestamp::now().as_secs(); - let created_at = if now > last_created_at { - now - } else { - last_created_at + 1 - }; - - let event = RepoState::build(repo_id, refs, head) - .custom_created_at(Timestamp::from_secs(created_at)) - .finalize_async(signer) - .await - .map_err(|e| format!("could not sign the state event: {e}"))?; - - Ok((event, created_at)) -} - -/// Ensure the relay is known and connected, then publish `event` to it. -/// -/// `Ok` only when the relay confirmed the event. -/// On a grasp relay the accept parks the event in purgatory, -/// which authorizes the paired git push. -async fn stage_event_on_relay( - client: &Client, - relay: &RelayUrl, - event: &Event, -) -> Result<(), String> { - client - .add_relay(relay) - .and_connect() - .await - .map_err(|e| format!("could not add relay {relay}: {e}"))?; - - let output = client - .send_event(event) - .to([relay.clone()]) - .await - .map_err(|e| format!("could not send the state event to {relay}: {e}"))?; - - if output.success.contains_key(relay) { - Ok(()) - } else { - let reason = output - .failed - .get(relay) - .cloned() - .unwrap_or_else(|| "relay did not confirm the event".to_owned()); - Err(reason) - } -} - -#[allow(clippy::too_many_arguments)] -async fn push_staged_to_grasps( - client: &Client, - signer: &UniversalSigner, - repo_id: &str, - refs: &[(String, String)], - head: Option<&str>, - path: &Path, - owner: &str, - servers: &[RelayUrl], - executor: &BackgroundExecutor, - push: impl Fn(&Path, &str, &str, &str) -> Result<(), Error>, -) -> PushOutcome { - let mut outcome = PushOutcome::default(); - - for relay in servers { - let Some(base) = grasp_base_url(relay) else { - outcome - .servers - .push(GraspServerResult::failed(relay.clone(), "no domain")); - continue; - }; - let git_url = format!("{base}/{owner}/{repo_id}.git"); - - let mut reason = None; - let mut last_created_at = 0; - // The last state event staged on this server, for the convergence - // probe below when every push attempt lost the stale-ref race. - let mut staged_event = None; - - 'server: for attempt in 1..=GRASP_PUSH_ATTEMPTS { - if attempt > 1 { - // Give the server's ingest a moment before re-staging. - executor.timer(GRASP_RETRY_DELAY).await; - } - - let (event, created_at) = - match sign_state_event(signer, repo_id, refs, head, last_created_at).await { - Ok(signed) => signed, - Err(e) => { - reason = Some(e); - break 'server; - } - }; - - last_created_at = created_at; - - // Stage the state event on this server's own relay. - // A failed stage means the grasp never parked the state, - // so the git push would be denied anyway: skip it (the eligibility gate). - if let Err(e) = stage_event_on_relay(client, relay, &event).await { - // One retry absorbs a relay connect blip, on the first - // attempt only. - if attempt == 1 && stage_event_on_relay(client, relay, &event).await.is_ok() { - // staged on the retry - } else { - reason = Some(e); - break 'server; - } - } - staged_event = Some(event.clone()); - - match push(path, &base, owner, repo_id) { - Ok(()) => { - keep_newest(&mut outcome.state_event, event); - break 'server; - } - Err(e) => { - let text = e.to_string(); - if attempt < GRASP_PUSH_ATTEMPTS && is_transient_grasp_denial(&text) { - reason = Some(text); - continue 'server; - } - reason = Some(text); - break 'server; } + Some(_) => continue, + None => return None, } } - - // The grasp's own background sync aligns refs to staged state - // events as soon as the objects land, which can beat every push - // attempt's compare-and-swap (`cannot lock ref ... but expected`). - // When the last denial was that race the sync has usually finished - // by now: verify the advertised refs and accept the server when the - // pushed data is already there. - if let Some(last_reason) = &reason - && is_stale_advertisement_race(last_reason) - && Repo::open(path) - .and_then(|repo| repo.remote_has_refs(&git_url, refs)) - .unwrap_or(false) - { - if let Some(event) = staged_event { - keep_newest(&mut outcome.state_event, event); - } - reason = None; - } - - match reason { - Some(reason) => { - log::warn!("grasp push failed: {relay}: {reason}"); - outcome - .servers - .push(GraspServerResult::failed(relay.clone(), reason)); - } - None => outcome.servers.push(GraspServerResult::ok(relay.clone())), - } } - - outcome } /// Split a stored bunker credential into the plain URI and the session key. @@ -1695,138 +1179,3 @@ fn extract_master_key(credential: &str) -> (&str, Keys) { None => (credential, Keys::generate()), } } - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn grasp_base_url_maps_schemes_like_ngit() { - let wss = RelayUrl::parse("wss://relay.ngit.dev").expect("url"); - assert_eq!( - grasp_base_url(&wss).as_deref(), - Some("https://relay.ngit.dev") - ); - - let ws = RelayUrl::parse("ws://localhost:8080").expect("url"); - assert_eq!( - grasp_base_url(&ws).as_deref(), - Some("http://localhost:8080") - ); - } - - #[test] - fn grasp_clone_url_matches_ngit_format() { - let relay = RelayUrl::parse("wss://gitnostr.com").expect("url"); - let url = grasp_clone_url(&relay, "npub1test", "my-repo").expect("url"); - assert_eq!( - url.to_string(), - "https://gitnostr.com/npub1test/my-repo.git" - ); - } - - #[test] - fn grasp06_prs_url_matches_ngit_format() { - assert_eq!( - grasp06_prs_url("https://relay.ngit.dev", "npub1author", "my-repo"), - "https://relay.ngit.dev/prs/npub1author/my-repo.git" - ); - // `ws://` grasp servers, local dev, keep their plain-HTTP base. - assert_eq!( - grasp06_prs_url("http://localhost:8080", "npub1author", "my-repo"), - "http://localhost:8080/prs/npub1author/my-repo.git" - ); - } - - #[test] - fn pr_clone_urls_orders_author_first_and_deduplicates() { - let prs = vec![ - Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"), - Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"), - ]; - let base = vec![ - Url::parse("https://a.example/npub1owner/repo.git").expect("url"), - Url::parse("https://b.example/npub1owner/repo.git").expect("url"), - Url::parse("https://b.example/npub1owner/repo.git").expect("url"), - ]; - - let urls = pr_clone_urls(prs, base); - assert_eq!( - urls.iter().map(ToString::to_string).collect::>(), - vec![ - "https://a.example/prs/npub1me/repo.git", - "https://a.example/npub1owner/repo.git", - "https://b.example/npub1owner/repo.git", - ] - ); - } - - #[test] - fn transient_grasp_denials_are_classified() { - // The exact server rejection that started this work: the state event - // had not reached the grasp's purgatory before the git push. - let reported = "remote: ERR authorisation failed: No state events in purgatory\n\ - fatal: the remote end hung up unexpectedly\n\ - error: failed to push some refs to 'https://relay.ngit.dev/...git'"; - assert!(is_transient_grasp_denial(reported)); - - // The other purgatory states a fresh event resolves. - assert!(is_transient_grasp_denial( - "remote: ERR authorisation failed: No matching state event found in purgatory" - )); - assert!(is_transient_grasp_denial( - "remote: ERR authorisation failed: 1 state event in purgatory from authorized \ - publisher but doesn't match push" - )); - assert!(is_transient_grasp_denial( - "remote: ERR authorisation failed: 2 state events in purgatory but none from \ - authorized publishers" - )); - assert!(is_transient_grasp_denial( - "remote: ERR authorisation failed: No repository announcement found" - )); - - // Rejections a fresh state event cannot fix are not retried. - assert!(!is_transient_grasp_denial( - "remote: ERR authorisation failed: not a maintainer of this repository" - )); - assert!(!is_transient_grasp_denial( - "fatal: unable to access 'https://relay.ngit.dev/...': The requested URL returned \ - error: 403" - )); - assert!(!is_transient_grasp_denial( - "fatal: unable to access 'https://relay.ngit.dev/...': Could not resolve host" - )); - } - - #[test] - fn stale_ref_races_are_retried() { - // The grasp's background sync aligned the ref to a parked state event - // between this push's advertisement and its ref transaction. The ref - // is usually already where the push wants it, so a retry converges. - let reported = "remote: error: cannot lock ref 'refs/heads/main': is at \ - cac2ac91b6f5fb8dfcb6962785babc6e65350cb3 but expected \ - bc5e892aa84dc6240a5fbcd59367a4857d26f49b\n\ - To https://relay.ngit.dev/npub1owner/signed-test.git\n\ - ! [remote rejected] main -> main (incorrect old value provided)\n\ - error: failed to push some refs to 'https://relay.ngit.dev/npub1owner/signed-test.git'"; - assert!(is_transient_grasp_denial(reported)); - assert!(is_stale_advertisement_race(reported)); - - // Markers match independently of the surrounding git output. - assert!(is_stale_advertisement_race( - "cannot lock ref 'refs/heads/main'" - )); - assert!(is_stale_advertisement_race( - "! [remote rejected] main -> main (incorrect old value provided)" - )); - - // A purgatory denial is not a stale-advertisement race. - assert!(!is_stale_advertisement_race("No state events in purgatory")); - - // A real divergence is a different error and stays permanent. - assert!(!is_transient_grasp_denial( - " ! [rejected] main -> main (non-fast-forward)" - )); - } -} diff --git a/crates/signed_state/src/bootstrap.rs b/crates/signed_state/src/bootstrap.rs new file mode 100644 index 0000000..906ac84 --- /dev/null +++ b/crates/signed_state/src/bootstrap.rs @@ -0,0 +1,99 @@ +use std::collections::HashMap; +use std::time::Duration; + +use anyhow::Error; +use nostr_connect::prelude::*; +use nostr_sdk::client::SyncSummary; +use nostr_sdk::prelude::*; +use signed_core::Filters; + +/// 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"]; + +/// 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))); + + let target: HashMap<&str, Vec> = BOOTSTRAP_RELAYS + .iter() + .map(|relay| (*relay, filters.clone())) + .collect(); + + client.subscribe(target).close_on(opts).await?; + + Ok(()) +} + +pub(crate) async fn sync_bootstrap_only( + client: &Client, + filter: Filter, + opts: SyncOptions, +) -> Result { + ensure_bootstrap_relays(client).await?; + + let output = client + .sync(filter) + .with(BOOTSTRAP_RELAYS) + .opts(opts) + .await?; + + Ok(output.value) +} + +/// The `g` tag servers of one kind-10317 grasp list event, in tag order. +fn grasp_list_servers(event: &Event) -> Vec { + event + .tags + .iter() + .filter(|tag| tag.kind() == "g") + .filter_map(|tag| tag.content()) + .filter_map(|url| RelayUrl::parse(url).ok()) + .collect() +} + +/// Grasp servers of the newest kind-10317 grasp list among `events`. +fn latest_grasp_list_servers(events: Vec) -> Vec { + events + .into_iter() + .max_by_key(|event| event.created_at) + .map(|event| grasp_list_servers(&event)) + .unwrap_or_default() +} + +pub async fn user_grasp_list_servers( + client: &Client, + user: PublicKey, +) -> Result, Error> { + let events: Vec = client + .database() + .query(Filters::grasp_list(user)) + .await? + .into_iter() + .collect(); + Ok(latest_grasp_list_servers(events)) +} diff --git a/crates/signed_state/src/checkouts.rs b/crates/signed_state/src/checkouts.rs index 0809319..d102197 100644 --- a/crates/signed_state/src/checkouts.rs +++ b/crates/signed_state/src/checkouts.rs @@ -11,7 +11,7 @@ use signed_git::Repo; use utils::same_repo_url; use crate::backend::{Backend, BackendEvent}; -use crate::git_store::repo_mirror_root; +use crate::git_store::Mirrors; use crate::local_repos::LocalReposStore; use crate::refresh::{RefreshGate, RefreshRequest}; use crate::repos::RepoListStore; @@ -308,7 +308,7 @@ impl CheckoutsStore { let announcements = RepoListStore::global(cx).read(cx).announcements.clone(); let scanned = LocalReposStore::global(cx).read(cx).repos.clone(); - let cache_root = repo_mirror_root().canonicalize().ok(); + let cache_root = Mirrors::root().canonicalize().ok(); let requested: Vec<(RepoAddr, Option)> = self .status_requested @@ -350,7 +350,8 @@ impl CheckoutsStore { facts.push((path.clone(), origin, root)); } - let associations = resolve_associations(&remembered, &facts, announcements.iter()); + let associations = + CheckoutsStore::resolve_associations(&remembered, &facts, announcements.iter()); // Missing directories are stale records, drop them. let associations: HashMap> = associations @@ -359,7 +360,7 @@ impl CheckoutsStore { .collect(); let (statuses, push_statuses) = - compute_statuses(&associations, &requested, &push_requested, true); + CheckoutsStore::compute_statuses(&associations, &requested, &push_requested, true); Ok::<_, Error>((associations, statuses, push_statuses)) }); @@ -485,7 +486,7 @@ impl CheckoutsStore { let work = cx.background_spawn(async move { let (statuses, push_statuses) = - compute_statuses(&associations, &requested, &push_requested, false); + CheckoutsStore::compute_statuses(&associations, &requested, &push_requested, false); Ok::<_, Error>((statuses, push_statuses)) }); @@ -519,190 +520,192 @@ impl CheckoutsStore { } } -fn resolve_associations<'a>( - remembered: &[Remembered], - scanned: &[(PathBuf, Option, Option)], - announcements: impl IntoIterator, -) -> HashMap> { - let announcements: Vec<&Announcement> = announcements.into_iter().collect(); - let mut out: HashMap> = HashMap::new(); +impl CheckoutsStore { + fn resolve_associations<'a>( + remembered: &[Remembered], + scanned: &[(PathBuf, Option, Option)], + announcements: impl IntoIterator, + ) -> HashMap> { + let announcements: Vec<&Announcement> = announcements.into_iter().collect(); + let mut out: HashMap> = HashMap::new(); - let mut sorted: Vec<&Remembered> = remembered.iter().collect(); - sorted.sort_by_key(|record| std::cmp::Reverse(record.last_used)); - for record in sorted { - let paths = out.entry(record.addr.clone()).or_default(); - if !paths.contains(&record.path) { - paths.push(record.path.clone()); + let mut sorted: Vec<&Remembered> = remembered.iter().collect(); + sorted.sort_by_key(|record| std::cmp::Reverse(record.last_used)); + for record in sorted { + let paths = out.entry(record.addr.clone()).or_default(); + if !paths.contains(&record.path) { + paths.push(record.path.clone()); + } } - } - for (path, origin, root) in scanned { - for announcement in &announcements { - let url_match = origin.as_deref().is_some_and(|origin| { - announcement - .clone - .iter() - .any(|url| same_repo_url(origin, url.as_str())) - }); + for (path, origin, root) in scanned { + for announcement in &announcements { + let url_match = origin.as_deref().is_some_and(|origin| { + announcement + .clone + .iter() + .any(|url| same_repo_url(origin, url.as_str())) + }); - let euc_match = root - .as_deref() - .is_some_and(|root| announcement.euc.as_deref() == Some(root)); + let euc_match = root + .as_deref() + .is_some_and(|root| announcement.euc.as_deref() == Some(root)); - if url_match || euc_match { - let paths = out.entry(announcement.addr()).or_default(); - if !paths.contains(path) { - paths.push(path.clone()); + if url_match || euc_match { + let paths = out.entry(announcement.addr()).or_default(); + if !paths.contains(path) { + paths.push(path.clone()); + } } } } + + out } - out -} + fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option { + let repo = Repo::try_open(path)?; + let branches = repo.branches().ok()?; -fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option { - let repo = Repo::try_open(path)?; - let branches = repo.branches().ok()?; + if branches.is_empty() || repo.is_dirty() { + return None; + } - if branches.is_empty() || repo.is_dirty() { - return None; + let branch = repo.current_branch()?; + let head = repo.head()?; + let base = announced_head + .filter(|name| branches.iter().any(|b| b == name)) + .map(str::to_owned) + .or_else(|| branches.iter().find(|b| *b == "main").cloned()) + .or_else(|| branches.first().cloned())?; + + if base == branch { + return None; + } + + let ahead = repo.commits_ahead(&base, &branch); + (ahead > 0).then_some(CheckoutStatus { + path: path.to_path_buf(), + branch, + head, + base, + ahead, + }) } - let branch = repo.current_branch()?; - let head = repo.head()?; - let base = announced_head - .filter(|name| branches.iter().any(|b| b == name)) - .map(str::to_owned) - .or_else(|| branches.iter().find(|b| *b == "main").cloned()) - .or_else(|| branches.first().cloned())?; + /// The `ready to push` status of one checkout of the user's own repository. + fn checkout_push_status(path: &Path, fetch: bool) -> Option { + let repo = Repo::try_open(path)?; + if repo.is_dirty() { + return None; + } - if base == branch { - return None; - } + let branch = repo.current_branch()?; + let head = repo.head()?; + let origin = repo.origin_url().ok().flatten()?; - let ahead = repo.commits_ahead(&base, &branch); - (ahead > 0).then_some(CheckoutStatus { - path: path.to_path_buf(), - branch, - head, - base, - ahead, - }) -} + if fetch { + repo.fetch_refs(&[origin], "+refs/heads/*:refs/remotes/origin/*") + .ok(); + } -/// The `ready to push` status of one checkout of the user's own repository. -fn checkout_push_status(path: &Path, fetch: bool) -> Option { - let repo = Repo::try_open(path)?; - if repo.is_dirty() { - return None; - } + let remote = format!("refs/remotes/origin/{branch}"); - let branch = repo.current_branch()?; - let head = repo.head()?; - let origin = repo.origin_url().ok().flatten()?; - - if fetch { - repo.fetch_refs(&[origin], "+refs/heads/*:refs/remotes/origin/*") - .ok(); - } - - let remote = format!("refs/remotes/origin/{branch}"); - - // A branch never fetched or pushed yet compares against the remote HEAD. - // The remote HEAD is the fork point in practice. - let base = if repo.ref_exists(&remote) { - remote - } else if repo.ref_exists("refs/remotes/origin/HEAD") { - "refs/remotes/origin/HEAD".to_owned() - } else { - return None; - }; - - let ahead = repo.commits_ahead(&base, &branch); - - (ahead > 0).then_some(CheckoutStatus { - path: path.to_path_buf(), - branch, - head, - base, - ahead, - }) -} - -/// Compute the requested statuses against the checkout paths of `associations`. -fn compute_statuses( - associations: &HashMap>, - requested: &[(RepoAddr, Option)], - push_requested: &[RepoAddr], - fetch: bool, -) -> ( - HashMap>, - HashMap>, -) { - let mut statuses: HashMap> = HashMap::new(); - for (addr, announced_head) in requested { - let Some(paths) = associations.get(addr) else { - continue; + // A branch never fetched or pushed yet compares against the remote HEAD. + // The remote HEAD is the fork point in practice. + let base = if repo.ref_exists(&remote) { + remote + } else if repo.ref_exists("refs/remotes/origin/HEAD") { + "refs/remotes/origin/HEAD".to_owned() + } else { + return None; }; - let list: Vec = paths - .iter() - .take(MAX_STATUS_CHECKOUTS) - .filter_map(|path| checkout_status(path, announced_head.as_deref())) - .collect(); + let ahead = repo.commits_ahead(&base, &branch); - if !list.is_empty() { - statuses.insert(addr.clone(), list); + (ahead > 0).then_some(CheckoutStatus { + path: path.to_path_buf(), + branch, + head, + base, + ahead, + }) + } + + /// Compute the requested statuses against the checkout paths of `associations`. + fn compute_statuses( + associations: &HashMap>, + requested: &[(RepoAddr, Option)], + push_requested: &[RepoAddr], + fetch: bool, + ) -> ( + HashMap>, + HashMap>, + ) { + let mut statuses: HashMap> = HashMap::new(); + for (addr, announced_head) in requested { + let Some(paths) = associations.get(addr) else { + continue; + }; + + let list: Vec = paths + .iter() + .take(MAX_STATUS_CHECKOUTS) + .filter_map(|path| Self::checkout_status(path, announced_head.as_deref())) + .collect(); + + if !list.is_empty() { + statuses.insert(addr.clone(), list); + } } - } - let mut push_statuses: HashMap> = HashMap::new(); - for addr in push_requested { - let Some(paths) = associations.get(addr) else { - continue; - }; + let mut push_statuses: HashMap> = HashMap::new(); + for addr in push_requested { + let Some(paths) = associations.get(addr) else { + continue; + }; - let list: Vec = paths - .iter() - .take(MAX_STATUS_CHECKOUTS) - .filter_map(|path| checkout_push_status(path, fetch)) - .collect(); + let list: Vec = paths + .iter() + .take(MAX_STATUS_CHECKOUTS) + .filter_map(|path| Self::checkout_push_status(path, fetch)) + .collect(); - if !list.is_empty() { - push_statuses.insert(addr.clone(), list); + if !list.is_empty() { + push_statuses.insert(addr.clone(), list); + } } + + (statuses, push_statuses) } - (statuses, push_statuses) -} + pub fn pr_proposes_checkout( + pr: &Event, + open: bool, + user: PublicKey, + checkout: &CheckoutStatus, + ) -> bool { + if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user { + return false; + } -pub fn pr_proposes_checkout( - pr: &Event, - open: bool, - user: PublicKey, - checkout: &CheckoutStatus, -) -> bool { - 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 + .iter() + .find(|t| t.kind() == "c") + .and_then(|t| t.content()) + .is_some_and(|tip| tip == checkout.head); + + branch_matches || tip_matches } - - 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 - .iter() - .find(|t| t.kind() == "c") - .and_then(|t| t.content()) - .is_some_and(|tip| tip == checkout.head); - - branch_matches || tip_matches } #[cfg(test)] @@ -763,7 +766,7 @@ mod tests { run(&["checkout", "-b", "feature"]); std::fs::write(path.join("feature.txt"), "x\n").expect("write"); commit("feature work"); - let status = checkout_status(&path, Some("main")).expect("status"); + let status = CheckoutsStore::checkout_status(&path, Some("main")).expect("status"); assert_eq!(status.branch, "feature"); assert_eq!(status.base, "main"); assert_eq!(status.ahead, 1); @@ -771,12 +774,12 @@ mod tests { // Dirty worktrees are never suggested. std::fs::write(path.join("uncommitted.txt"), "y\n").expect("write"); - assert!(checkout_status(&path, Some("main")).is_none()); + assert!(CheckoutsStore::checkout_status(&path, Some("main")).is_none()); run(&["checkout", "--", "."]); // Even on main, nothing to propose. run(&["checkout", "main"]); - assert_eq!(checkout_status(&path, Some("main")), None); + assert_eq!(CheckoutsStore::checkout_status(&path, Some("main")), None); } #[test] @@ -822,13 +825,13 @@ mod tests { }; // A fresh clone has nothing to push. - assert_eq!(checkout_push_status(&checkout, true), None); + assert_eq!(CheckoutsStore::checkout_push_status(&checkout, true), None); // One local commit, ready to push, counted against the remote. std::fs::write(checkout.join("work.txt"), "x\n").expect("write"); run(&["add", "-A"]); run(&["commit", "-m", "local work"]); - let status = checkout_push_status(&checkout, true).expect("status"); + let status = CheckoutsStore::checkout_push_status(&checkout, true).expect("status"); assert_eq!(status.branch, "main"); assert_eq!(status.base, "refs/remotes/origin/main"); assert_eq!(status.ahead, 1); @@ -836,12 +839,12 @@ mod tests { // The local-only pass reads the tracking refs, no fetch needed: // a commit lands locally long before the remote is reconciled. - let local = checkout_push_status(&checkout, false).expect("local status"); + let local = CheckoutsStore::checkout_push_status(&checkout, false).expect("local status"); assert_eq!(local.ahead, 1); // After the push the same commit is on the remote, idle again. run(&["push", "origin", "main"]); - assert_eq!(checkout_push_status(&checkout, true), None); + assert_eq!(CheckoutsStore::checkout_push_status(&checkout, true), None); // A commit made by someone else on the remote must not count as local work. // It is behind, not ahead. @@ -860,6 +863,6 @@ mod tests { std::fs::write(remote.join("other.txt"), "y\n").expect("write"); remote_run(&["add", "-A"]); remote_run(&["commit", "-m", "remote work"]); - assert_eq!(checkout_push_status(&checkout, true), None); + assert_eq!(CheckoutsStore::checkout_push_status(&checkout, true), None); } } diff --git a/crates/signed_state/src/git_store.rs b/crates/signed_state/src/git_store.rs index 961136c..c4072a2 100644 --- a/crates/signed_state/src/git_store.rs +++ b/crates/signed_state/src/git_store.rs @@ -8,34 +8,40 @@ use signed_git::{GitCache, Repo}; static GIT_CACHE: OnceLock = OnceLock::new(); -fn git_cache() -> &'static GitCache { - GIT_CACHE - .get() - .expect("git cache is initialized by signed_state::init") -} +/// The global git repository mirror cache. +pub struct Mirrors; -/// The root directory of the repository mirrors. -pub(crate) fn repo_mirror_root() -> PathBuf { - git_cache().root().to_path_buf() -} +impl Mirrors { + /// Install the global mirror cache root, once. + pub fn install(root: impl Into) { + if GIT_CACHE.set(GitCache::new(root.into())).is_err() { + log::warn!("git cache root is already set, keeping the first one"); + } + } -/// The on-disk path of the mirror of `addr`. -pub fn repo_mirror_path(addr: &RepoAddr) -> PathBuf { - git_cache().repo_path(addr) -} + fn cache() -> &'static GitCache { + GIT_CACHE + .get() + .expect("git cache is initialized by signed_state::init") + } -/// Open the mirror of `addr`, if it has been cloned. -pub fn open_repo_mirror(addr: &RepoAddr) -> Result> { - git_cache().open(addr) -} + /// The root directory of the repository mirrors. + pub(crate) fn root() -> PathBuf { + Self::cache().root().to_path_buf() + } -/// Open the mirror of `addr`, cloning it first when it does not exist yet. -pub fn ensure_repo_mirror(addr: &RepoAddr, clone_urls: &[Url]) -> Result { - git_cache().ensure_clone(addr, clone_urls) -} + /// The on-disk path of the mirror of `addr`. + pub fn path(addr: &RepoAddr) -> PathBuf { + Self::cache().repo_path(addr) + } -pub(crate) fn set_git_cache(root: impl Into) { - if GIT_CACHE.set(GitCache::new(root.into())).is_err() { - log::warn!("git cache root is already set, keeping the first one"); + /// Open the mirror of `addr`, if it has been cloned. + pub fn open(addr: &RepoAddr) -> Result> { + Self::cache().open(addr) + } + + /// Open the mirror of `addr`, cloning it first when it does not exist yet. + pub fn ensure(addr: &RepoAddr, clone_urls: &[Url]) -> Result { + Self::cache().ensure_clone(addr, clone_urls) } } diff --git a/crates/signed_state/src/inbox.rs b/crates/signed_state/src/inbox.rs index d7be438..96e56b7 100644 --- a/crates/signed_state/src/inbox.rs +++ b/crates/signed_state/src/inbox.rs @@ -36,7 +36,7 @@ impl Inbox { cx.notify(); let backend = Backend::global(cx); - let work = cx.background_spawn(async move { load_state(&client, me).await }); + let work = cx.background_spawn(async move { Self::load_state(&client, me).await }); cx.spawn(async move |this, cx| { let loaded = work.await; @@ -82,7 +82,7 @@ impl Inbox { let state = self.state.clone(); let task: Task> = cx.background_spawn(async move { - if let Err(error) = save_state(&client, me, &state).await { + if let Err(error) = Self::save_state(&client, me, &state).await { log::warn!("failed to save inbox state: {error}"); } Ok(()) @@ -90,135 +90,147 @@ impl Inbox { task.detach(); } -} -pub async fn query_inbox( - client: &Client, - me: PublicKey, - state: &InboxReadState, -) -> Result<(Vec, usize), Error> { - let deletion_events = client.database().query(Filters::deletions()).await?; - let deletions = Deletions::from_events(deletion_events); + /// The notifications and authored activity of `me`, grouped into inbox items. + pub async fn query( + client: &Client, + me: PublicKey, + state: &InboxReadState, + ) -> Result<(Vec, usize), Error> { + let deletion_events = client.database().query(Filters::deletions()).await?; + let deletions = Deletions::from_events(deletion_events); - let (notification_events, mut by_id) = fetch_notifications(client, me, &deletions).await?; + let (notification_events, mut by_id) = + Self::fetch_notifications(client, me, &deletions).await?; - let mut activity = Vec::new(); - for event in client - .database() - .query(Filters::authored_activity(me)) - .await? - { - if deletions.is_deleted(&event) || !event.is_git_activity() { - continue; - } - by_id.entry(event.id).or_insert_with(|| event.clone()); - activity.push(event); - } - - let items = inbox::group(notification_events, activity, me, state, &|id| { - by_id.get(&id).cloned() - }); - - let unread_count = items.iter().filter(|item| item.is_unread()).count(); - - Ok((items, unread_count)) -} - -/// `d` tag identifying the inbox state event of `me`. -fn inbox_state_d_tag(me: PublicKey) -> String { - format!("signed-inbox-state:{}", me.to_hex()) -} - -async fn load_state(client: &Client, me: PublicKey) -> Result, Error> { - let filter = Filter::new() - .kind(Kind::ApplicationSpecificData) - .identifier(inbox_state_d_tag(me)); - - let events = client.database().query(filter).await?; - - let Some(event) = events.into_iter().max_by_key(|event| event.created_at) else { - return Ok(None); - }; - - match serde_json::from_str(&event.content) { - Ok(state) => Ok(Some(state)), - Err(error) => { - log::warn!("ignoring unreadable inbox state {}: {error}", event.id); - Ok(None) - } - } -} - -/// Sign with a random key and store locally. -async fn save_state(client: &Client, me: PublicKey, state: &InboxReadState) -> Result<(), Error> { - let event = EventBuilder::new(Kind::ApplicationSpecificData, serde_json::to_string(state)?) - .tags([Tag::identifier(inbox_state_d_tag(me))]) - .finalize(&Keys::generate())?; - - client.database().save_event(&event).await?; - - Ok(()) -} - -/// Notification events and a lookup of every ancestor they reference. -async fn fetch_notifications( - client: &Client, - me: PublicKey, - deletions: &Deletions, -) -> Result<(Vec, HashMap), Error> { - let mut notifications: Vec = Vec::new(); - let mut by_id: HashMap = HashMap::new(); - - for filter in Filters::notifications(me) { - for event in client.database().query(filter).await? { - if deletions.is_deleted(&event) { - continue; - } - - if by_id.insert(event.id, event.clone()).is_none() { - notifications.push(event); - } - } - } - - let mut pending: Vec = notifications.iter().flat_map(event_references).collect(); - let mut seen: HashSet = by_id.keys().copied().collect(); - - loop { - pending.retain(|id| seen.insert(*id)); - - if pending.is_empty() { - break; - } - - let ancestors = client + let mut activity = Vec::new(); + for event in client .database() - .query(Filter::new().ids(pending.iter().copied())) - .await?; - - let mut next = Vec::new(); - - for event in ancestors { - if deletions.is_deleted(&event) { + .query(Filters::authored_activity(me)) + .await? + { + if deletions.is_deleted(&event) || !event.is_git_activity() { continue; } - next.extend(event_references(&event).filter(|id| !seen.contains(id))); - by_id.entry(event.id).or_insert(event); + by_id.entry(event.id).or_insert_with(|| event.clone()); + activity.push(event); } - pending = next; + let items = inbox::group(notification_events, activity, me, state, &|id| { + by_id.get(&id).cloned() + }); + + let unread_count = items.iter().filter(|item| item.is_unread()).count(); + + Ok((items, unread_count)) + } +} + +impl Inbox { + /// `d` tag identifying the inbox state event of `me`. + fn inbox_state_d_tag(me: PublicKey) -> String { + format!("signed-inbox-state:{}", me.to_hex()) } - Ok((notifications, by_id)) -} + async fn load_state(client: &Client, me: PublicKey) -> Result, Error> { + let filter = Filter::new() + .kind(Kind::ApplicationSpecificData) + .identifier(Self::inbox_state_d_tag(me)); -/// Event ids referenced by `event` through its `e` and `E` tags. -fn event_references(event: &Event) -> impl Iterator + '_ { - event.tags.iter().filter_map(|tag| { - if tag.kind() != "e" && tag.kind() != "E" { - return None; + let events = client.database().query(filter).await?; + + let Some(event) = events.into_iter().max_by_key(|event| event.created_at) else { + return Ok(None); + }; + + match serde_json::from_str(&event.content) { + Ok(state) => Ok(Some(state)), + Err(error) => { + log::warn!("ignoring unreadable inbox state {}: {error}", event.id); + Ok(None) + } } - tag.content() - .and_then(|content| EventId::from_hex(content).ok()) - }) + } + + /// Sign with a random key and store locally. + async fn save_state( + client: &Client, + me: PublicKey, + state: &InboxReadState, + ) -> Result<(), Error> { + let event = EventBuilder::new(Kind::ApplicationSpecificData, serde_json::to_string(state)?) + .tags([Tag::identifier(Self::inbox_state_d_tag(me))]) + .finalize(&Keys::generate())?; + + client.database().save_event(&event).await?; + + Ok(()) + } + + /// Notification events and a lookup of every ancestor they reference. + async fn fetch_notifications( + client: &Client, + me: PublicKey, + deletions: &Deletions, + ) -> Result<(Vec, HashMap), Error> { + let mut notifications: Vec = Vec::new(); + let mut by_id: HashMap = HashMap::new(); + + for filter in Filters::notifications(me) { + for event in client.database().query(filter).await? { + if deletions.is_deleted(&event) { + continue; + } + + if by_id.insert(event.id, event.clone()).is_none() { + notifications.push(event); + } + } + } + + let mut pending: Vec = notifications + .iter() + .flat_map(Self::event_references) + .collect(); + + let mut seen: HashSet = by_id.keys().copied().collect(); + + loop { + pending.retain(|id| seen.insert(*id)); + + if pending.is_empty() { + break; + } + + let ancestors = client + .database() + .query(Filter::new().ids(pending.iter().copied())) + .await?; + + let mut next = Vec::new(); + + for event in ancestors { + if deletions.is_deleted(&event) { + continue; + } + next.extend(Self::event_references(&event).filter(|id| !seen.contains(id))); + by_id.entry(event.id).or_insert(event); + } + + pending = next; + } + + Ok((notifications, by_id)) + } + + /// Event ids referenced by `event` through its `e` and `E` tags. + fn event_references(event: &Event) -> impl Iterator + '_ { + event.tags.iter().filter_map(|tag| { + if tag.kind() != "e" && tag.kind() != "E" { + return None; + } + tag.content() + .and_then(|content| EventId::from_hex(content).ok()) + }) + } } diff --git a/crates/signed_state/src/lib.rs b/crates/signed_state/src/lib.rs index 217cc56..163d5be 100644 --- a/crates/signed_state/src/lib.rs +++ b/crates/signed_state/src/lib.rs @@ -1,24 +1,27 @@ mod backend; +mod bootstrap; mod checkouts; mod git_store; mod inbox; mod local_repos; mod profile; +mod push; mod refresh; mod repo; mod repos; use std::path::{Path, PathBuf}; -pub use backend::{Backend, BackendEvent, user_grasp_list_servers}; -pub use checkouts::{CheckoutStatus, CheckoutsStore, pr_proposes_checkout}; -use git_store::set_git_cache; -pub use git_store::{ensure_repo_mirror, open_repo_mirror, repo_mirror_path}; +pub use backend::{Backend, BackendEvent}; +pub use bootstrap::user_grasp_list_servers; +pub use checkouts::{CheckoutStatus, CheckoutsStore}; +pub use git_store::Mirrors; use gpui::{App, AppContext}; -pub use inbox::{Inbox, query_inbox}; -pub use local_repos::{LocalReposStore, ResolvedLocalRepo, resolve_local_repos}; +pub use inbox::Inbox; +pub use local_repos::{LocalReposStore, ResolvedLocalRepo}; pub use nostr_sdk::prelude::Timestamp; pub use profile::{Profile, ProfileStore}; +pub use push::{GraspServer, PushOutcome}; pub use refresh::{RefreshGate, RefreshRequest}; pub use repo::RepoStore; pub use repos::RepoListStore; @@ -42,7 +45,7 @@ pub fn init( let (client, signer) = (backend.client, backend.signer); // Set Git cache for the repos root - set_git_cache(repos_root); + Mirrors::install(repos_root); // Set global stores for the backend Backend::set_global(cx.new(|cx| Backend::new(client, signer, cx)), cx); diff --git a/crates/signed_state/src/local_repos.rs b/crates/signed_state/src/local_repos.rs index c1ad4fb..12b84d0 100644 --- a/crates/signed_state/src/local_repos.rs +++ b/crates/signed_state/src/local_repos.rs @@ -137,40 +137,42 @@ impl ResolvedLocalRepo { } /// Resolve the scanned repositories against the known announcements. -pub fn resolve_local_repos( - repos: &[LocalRepo], - known: &[Announcement], - own: &[Announcement], -) -> Vec { - let shown: HashSet = own.iter().map(Announcement::addr).collect(); +impl LocalReposStore { + pub fn resolve( + repos: &[LocalRepo], + known: &[Announcement], + own: &[Announcement], + ) -> Vec { + let shown: HashSet = own.iter().map(Announcement::addr).collect(); - repos - .iter() - .filter_map(|repo| { - let addr = local_repo_addr(repo); + repos + .iter() + .filter_map(|repo| { + let addr = local_repo_addr(repo); - if let Some(addr) = &addr - && shown.contains(addr) - { - return None; - } + if let Some(addr) = &addr + && shown.contains(addr) + { + return None; + } - let announcement = addr - .as_ref() - .and_then(|addr| { - known - .iter() - .find(|announcement| announcement.addr() == *addr) + let announcement = addr + .as_ref() + .and_then(|addr| { + known + .iter() + .find(|announcement| announcement.addr() == *addr) + }) + .cloned(); + + Some(ResolvedLocalRepo { + path: repo.path.clone(), + nip34: repo.nip34.clone(), + announcement, }) - .cloned(); - - Some(ResolvedLocalRepo { - path: repo.path.clone(), - nip34: repo.nip34.clone(), - announcement, }) - }) - .collect() + .collect() + } } #[cfg(test)] @@ -222,7 +224,7 @@ mod tests { let repo = bound(KEY, "mine"); let own = std::slice::from_ref(&own); - assert!(resolve_local_repos(&[repo], own, own).is_empty()); + assert!(LocalReposStore::resolve(&[repo], own, own).is_empty()); } #[test] @@ -230,7 +232,7 @@ mod tests { let known = announcement(OTHER_KEY, "theirs"); let repo = bound(OTHER_KEY, "theirs"); - let resolved = resolve_local_repos(&[repo], std::slice::from_ref(&known), &[]); + let resolved = LocalReposStore::resolve(&[repo], std::slice::from_ref(&known), &[]); assert_eq!(resolved.len(), 1); assert_eq!(resolved[0].announcement.as_ref(), Some(&known)); diff --git a/crates/signed_state/src/profile.rs b/crates/signed_state/src/profile.rs index 02f84b5..daeffc3 100644 --- a/crates/signed_state/src/profile.rs +++ b/crates/signed_state/src/profile.rs @@ -10,7 +10,8 @@ use gpui::{ use nostr_sdk::prelude::*; use utils::shorten_pubkey; -use crate::backend::{Backend, BackendEvent, sync_bootstrap_only}; +use crate::backend::{Backend, BackendEvent}; +use crate::bootstrap::sync_bootstrap_only; /// How long to wait for more requests before firing a batched fetch. const BATCH_TIMEOUT: Duration = Duration::from_millis(500); diff --git a/crates/signed_state/src/push.rs b/crates/signed_state/src/push.rs new file mode 100644 index 0000000..9bd06af --- /dev/null +++ b/crates/signed_state/src/push.rs @@ -0,0 +1,577 @@ +use std::path::Path; +use std::time::Duration; + +use anyhow::{Error, bail}; +use gpui::BackgroundExecutor; +use nostr::event::IntoEventBuilder; +use nostr::prelude::Url; +use nostr_sdk::prelude::*; +use signed_core::RepoState; +use signed_git::Repo; +use signed_nostr::UniversalSigner; + +pub(crate) const GRASP_PUSH_ATTEMPTS: usize = 3; + +/// Pause before re-staging a state event after a transient denial. +const GRASP_RETRY_DELAY: Duration = Duration::from_secs(1); + +/// Base URL of a grasp server, `https://`. +/// +/// `ws://` grasp servers use `http://`, like ngit. +pub(crate) fn grasp_base_url(relay: &RelayUrl) -> Option { + // `domain()` drops the port. + let parsed = Url::parse(relay.as_str()).ok()?; + let host = parsed.host_str()?; + let port = parsed.port().map(|p| format!(":{p}")).unwrap_or_default(); + // `ws://` grasp servers, e.g. local dev relays, speak plain HTTP. + let scheme = if relay.scheme().is_secure() { + "https" + } else { + "http" + }; + Some(format!("{scheme}://{host}{port}")) +} + +pub(crate) fn grasp_clone_url(relay: &RelayUrl, owner: &str, repo_id: &str) -> Option { + let base = grasp_base_url(relay)?; + Url::parse(&format!("{base}/{owner}/{repo_id}.git")).ok() +} + +/// GRASP-06 contributor namespace URL of a pull request tip. +pub(crate) fn grasp06_prs_url(base_url: &str, npub: &str, repo_id: &str) -> String { + format!("{base_url}/prs/{npub}/{repo_id}.git") +} + +/// Assemble the `clone` URLs of a pull request. +/// +/// The author's GRASP-06 `/prs/` URLs come first. +pub(crate) fn pr_clone_urls(prs_urls: Vec, base_clone_urls: Vec) -> Vec { + let mut seen = std::collections::HashSet::new(); + let mut urls = Vec::new(); + for url in prs_urls.into_iter().chain(base_clone_urls) { + if seen.insert(url.to_string()) { + urls.push(url); + } + } + urls +} + +#[derive(Debug, Clone)] +pub struct GraspServer { + relay: RelayUrl, + /// `None` when the server accepted the data, the reason otherwise. + reason: Option, +} + +impl GraspServer { + fn ok(relay: RelayUrl) -> Self { + Self { + relay, + reason: None, + } + } + + fn failed(relay: RelayUrl, reason: impl Into) -> Self { + Self { + relay, + reason: Some(reason.into()), + } + } + + pub fn relay(&self) -> &RelayUrl { + &self.relay + } + + /// `None` when the server accepted the data, the reason otherwise. + pub fn reason(&self) -> Option<&str> { + self.reason.as_deref() + } +} + +#[derive(Debug, Clone, Default)] +pub struct PushOutcome { + /// Per-server results, in the order the servers were listed. + pub servers: Vec, + /// The newest state event a grasp relay accepted for this push, if any. + /// + /// Broadcast to the other relays once a git server holds the data. + pub state_event: Option, +} + +impl PushOutcome { + pub fn accepted(&self) -> usize { + self.servers + .iter() + .filter(|server| server.reason.is_none()) + .count() + } + + fn failing(&self) -> impl Iterator { + self.servers.iter().filter(|server| server.reason.is_some()) + } + + pub fn failure_summary(&self) -> String { + self.failing() + .map(|server| { + let reason = + utils::flatten_whitespace(server.reason.as_deref().unwrap_or("unknown error")); + format!("{}: {reason}", server.relay) + }) + .collect::>() + .join("; ") + } + + /// A warning for a push only some grasp servers accepted. + /// + /// `None` when every server accepted the push or nothing was pushed. + pub fn partial_warning(&self) -> Option { + let accepted = self.accepted(); + if self.servers.is_empty() || accepted == self.servers.len() { + return None; + } + Some(format!( + "Pushed to {accepted} of {} grasp servers: {}. Republish to sync.", + self.servers.len(), + self.failure_summary() + )) + } +} + +/// Reasons a push attempt should be retried with a freshly staged state +/// event and a fresh git advertisement. +/// +/// Two families are retried: +/// +/// - **Purgatory denials**: the grasp server sends these when the state +/// event for the push has not reached its purgatory yet. Re-staging a +/// fresh event resolves them. +/// - **Stale advertisement races**: `git receive-pack` compares each ref +/// update against the value it advertised when the push started. The grasp +/// server's own background sync can move a ref in between - typically by +/// aligning the repository to a parked state event once the objects of an +/// earlier attempt land - so the compare-and-swap fails with `cannot lock +/// ref` / `incorrect old value provided`. A retry against the fresh +/// advertisement converges, and when the race is lost the pushed data is +/// usually already on the server (see `is_stale_advertisement_race` and +/// the convergence probe in `push_staged_to_grasps`). +/// +/// Other rejections are not retried. +fn is_transient_grasp_denial(stderr: &str) -> bool { + let error = stderr.to_lowercase(); + [ + "no state events in purgatory", + "no matching state event", + "doesn't match push", + "none from authorized publishers", + "no repository announcement found", + "cannot lock ref", + "incorrect old value provided", + ] + .iter() + .any(|marker| error.contains(marker)) +} + +/// A push rejected because `git receive-pack`'s compare-and-swap lost to the +/// grasp server's own background ref alignment: the ref moved between this +/// push's advertisement and its ref transaction (`cannot lock ref ... is at +/// ... but expected ...` / `incorrect old value provided`). The pushed data +/// is usually already on the server by then. +fn is_stale_advertisement_race(stderr: &str) -> bool { + let error = stderr.to_lowercase(); + error.contains("cannot lock ref") || error.contains("incorrect old value provided") +} + +/// Keep `event` as the push's fan-out state event when it is newer than the +/// current one. All staged events carry the same refs; the newest timestamp +/// wins on the relays. +fn keep_newest(state_event: &mut Option, event: Event) { + if state_event + .as_ref() + .is_none_or(|current| event.created_at > current.created_at) + { + *state_event = Some(event); + } +} + +/// The grasp push pipeline: stage a signed state event on each server's +/// relay, then push the git data, retrying transient denials. +#[derive(Clone)] +pub(crate) struct GraspPush { + client: Client, + signer: UniversalSigner, +} + +impl GraspPush { + pub(crate) fn new(client: Client, signer: UniversalSigner) -> Self { + Self { client, signer } + } + + /// The event was accepted by at least one relay, or a descriptive error otherwise. + pub(crate) fn require_relay_accepted( + output: SendEventOutput, + event: Event, + ) -> Result { + if output.success.is_empty() && !output.failed.is_empty() { + let reasons = output + .failed + .values() + .cloned() + .collect::>() + .join(", "); + bail!("event not accepted by any relay: {reasons}"); + } + + Ok(event) + } + + /// Sign and broadcast `builder`, logging rather than surfacing failures. + /// + /// Used for best-effort identity bootstrap events, where a relay hiccup + /// should not block sign-up. + pub(crate) async fn publish_best_effort(&self, builder: EventBuilder) { + let result: Result<(), Error> = async { + let event = builder.finalize_async(&self.signer).await?; + let output = self.client.send_event(&event).broadcast().await?; + Self::require_relay_accepted(output, event)?; + Ok(()) + } + .await; + + if let Err(e) = result { + log::warn!("failed to publish identity bootstrap event: {e}"); + } + } + + /// Sign `builder`, broadcast the event and require a relay to accept it. + /// Returns the signed event. + pub(crate) async fn publish_one(&self, builder: EventBuilder) -> Result { + let event = builder.finalize_async(&self.signer).await?; + self.send_accepted(event).await + } + + /// Broadcast an already signed event and require a relay to accept it. + /// Returns the event. + pub(crate) async fn send_accepted(&self, event: Event) -> Result { + let output = self.client.send_event(&event).broadcast().await?; + Self::require_relay_accepted(output, event) + } + + /// Sign and send a single NIP-09 deletion request for `event`. + pub(crate) async fn retract_event(&self, event: &Event) -> Result<(), Error> { + let builder = EventDeletionRequest::new() + .id(event.id) + .into_event_builder(); + + let deletion = builder.finalize_async(&self.signer).await?; + self.client.send_event(&deletion).broadcast().await?; + + Ok(()) + } + + /// Sign a fresh kind `30618` state event for the push. + /// + /// `last_created_at` is the timestamp of the previous event signed for this push. + /// Retries within the same second get the next second: a grasp relay + /// treats a same-id resend as a duplicate and does not re-run its ingest, + /// so an identical resend cannot re-park a state event lost from its purgatory. + async fn sign_state_event( + &self, + repo_id: &str, + refs: &[(String, String)], + head: Option<&str>, + last_created_at: u64, + ) -> Result<(Event, u64), String> { + let now = Timestamp::now().as_secs(); + let created_at = if now > last_created_at { + now + } else { + last_created_at + 1 + }; + + let event = RepoState::build(repo_id, refs, head) + .custom_created_at(Timestamp::from_secs(created_at)) + .finalize_async(&self.signer) + .await + .map_err(|e| format!("could not sign the state event: {e}"))?; + + Ok((event, created_at)) + } + + /// Ensure the relay is known and connected, then publish `event` to it. + /// + /// `Ok` only when the relay confirmed the event. + /// On a grasp relay the accept parks the event in purgatory, + /// which authorizes the paired git push. + async fn stage_event_on_relay(&self, relay: &RelayUrl, event: &Event) -> Result<(), String> { + self.client + .add_relay(relay) + .and_connect() + .await + .map_err(|e| format!("could not add relay {relay}: {e}"))?; + + let output = self + .client + .send_event(event) + .to([relay.clone()]) + .await + .map_err(|e| format!("could not send the state event to {relay}: {e}"))?; + + if output.success.contains_key(relay) { + Ok(()) + } else { + let reason = output + .failed + .get(relay) + .cloned() + .unwrap_or_else(|| "relay did not confirm the event".to_owned()); + Err(reason) + } + } + + #[allow(clippy::too_many_arguments)] + pub(crate) async fn push_staged_to_grasps( + &self, + repo_id: &str, + refs: &[(String, String)], + head: Option<&str>, + path: &Path, + owner: &str, + servers: &[RelayUrl], + executor: &BackgroundExecutor, + push: impl Fn(&Path, &str, &str, &str) -> Result<(), Error>, + ) -> PushOutcome { + let mut outcome = PushOutcome::default(); + + for relay in servers { + let Some(base) = grasp_base_url(relay) else { + outcome + .servers + .push(GraspServer::failed(relay.clone(), "no domain")); + continue; + }; + let git_url = format!("{base}/{owner}/{repo_id}.git"); + + let mut reason = None; + let mut last_created_at = 0; + // The last state event staged on this server, for the convergence + // probe below when every push attempt lost the stale-ref race. + let mut staged_event = None; + + 'server: for attempt in 1..=GRASP_PUSH_ATTEMPTS { + if attempt > 1 { + // Give the server's ingest a moment before re-staging. + executor.timer(GRASP_RETRY_DELAY).await; + } + + let (event, created_at) = match self + .sign_state_event(repo_id, refs, head, last_created_at) + .await + { + Ok(signed) => signed, + Err(e) => { + reason = Some(e); + break 'server; + } + }; + + last_created_at = created_at; + + // Stage the state event on this server's own relay. + // A failed stage means the grasp never parked the state, + // so the git push would be denied anyway: skip it (the eligibility gate). + if let Err(e) = self.stage_event_on_relay(relay, &event).await { + // One retry absorbs a relay connect blip, on the first + // attempt only. + if attempt == 1 && self.stage_event_on_relay(relay, &event).await.is_ok() { + // staged on the retry + } else { + reason = Some(e); + break 'server; + } + } + staged_event = Some(event.clone()); + + match push(path, &base, owner, repo_id) { + Ok(()) => { + keep_newest(&mut outcome.state_event, event); + break 'server; + } + Err(e) => { + let text = e.to_string(); + if attempt < GRASP_PUSH_ATTEMPTS && is_transient_grasp_denial(&text) { + reason = Some(text); + continue 'server; + } + reason = Some(text); + break 'server; + } + } + } + + // The grasp's own background sync aligns refs to staged state + // events as soon as the objects land, which can beat every push + // attempt's compare-and-swap (`cannot lock ref ... but expected`). + // When the last denial was that race the sync has usually finished + // by now: verify the advertised refs and accept the server when the + // pushed data is already there. + if let Some(last_reason) = &reason + && is_stale_advertisement_race(last_reason) + && Repo::open(path) + .and_then(|repo| repo.remote_has_refs(&git_url, refs)) + .unwrap_or(false) + { + if let Some(event) = staged_event { + keep_newest(&mut outcome.state_event, event); + } + reason = None; + } + + match reason { + Some(reason) => { + log::warn!("grasp push failed: {relay}: {reason}"); + outcome + .servers + .push(GraspServer::failed(relay.clone(), reason)); + } + None => outcome.servers.push(GraspServer::ok(relay.clone())), + } + } + + outcome + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn grasp_base_url_maps_schemes_like_ngit() { + let wss = RelayUrl::parse("wss://relay.ngit.dev").expect("url"); + assert_eq!( + grasp_base_url(&wss).as_deref(), + Some("https://relay.ngit.dev") + ); + + let ws = RelayUrl::parse("ws://localhost:8080").expect("url"); + assert_eq!( + grasp_base_url(&ws).as_deref(), + Some("http://localhost:8080") + ); + } + + #[test] + fn grasp_clone_url_matches_ngit_format() { + let relay = RelayUrl::parse("wss://gitnostr.com").expect("url"); + let url = grasp_clone_url(&relay, "npub1test", "my-repo").expect("url"); + assert_eq!( + url.to_string(), + "https://gitnostr.com/npub1test/my-repo.git" + ); + } + + #[test] + fn grasp06_prs_url_matches_ngit_format() { + assert_eq!( + grasp06_prs_url("https://relay.ngit.dev", "npub1author", "my-repo"), + "https://relay.ngit.dev/prs/npub1author/my-repo.git" + ); + // `ws://` grasp servers, local dev, keep their plain-HTTP base. + assert_eq!( + grasp06_prs_url("http://localhost:8080", "npub1author", "my-repo"), + "http://localhost:8080/prs/npub1author/my-repo.git" + ); + } + + #[test] + fn pr_clone_urls_orders_author_first_and_deduplicates() { + let prs = vec![ + Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"), + Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"), + ]; + let base = vec![ + Url::parse("https://a.example/npub1owner/repo.git").expect("url"), + Url::parse("https://b.example/npub1owner/repo.git").expect("url"), + Url::parse("https://b.example/npub1owner/repo.git").expect("url"), + ]; + + let urls = pr_clone_urls(prs, base); + assert_eq!( + urls.iter().map(ToString::to_string).collect::>(), + vec![ + "https://a.example/prs/npub1me/repo.git", + "https://a.example/npub1owner/repo.git", + "https://b.example/npub1owner/repo.git", + ] + ); + } + + #[test] + fn transient_grasp_denials_are_classified() { + // The exact server rejection that started this work: the state event + // had not reached the grasp's purgatory before the git push. + let reported = "remote: ERR authorisation failed: No state events in purgatory\n\ + fatal: the remote end hung up unexpectedly\n\ + error: failed to push some refs to 'https://relay.ngit.dev/...git'"; + assert!(is_transient_grasp_denial(reported)); + + // The other purgatory states a fresh event resolves. + assert!(is_transient_grasp_denial( + "remote: ERR authorisation failed: No matching state event found in purgatory" + )); + assert!(is_transient_grasp_denial( + "remote: ERR authorisation failed: 1 state event in purgatory from authorized \ + publisher but doesn't match push" + )); + assert!(is_transient_grasp_denial( + "remote: ERR authorisation failed: 2 state events in purgatory but none from \ + authorized publishers" + )); + assert!(is_transient_grasp_denial( + "remote: ERR authorisation failed: No repository announcement found" + )); + + // Rejections a fresh state event cannot fix are not retried. + assert!(!is_transient_grasp_denial( + "remote: ERR authorisation failed: not a maintainer of this repository" + )); + assert!(!is_transient_grasp_denial( + "fatal: unable to access 'https://relay.ngit.dev/...': The requested URL returned \ + error: 403" + )); + assert!(!is_transient_grasp_denial( + "fatal: unable to access 'https://relay.ngit.dev/...': Could not resolve host" + )); + } + + #[test] + fn stale_ref_races_are_retried() { + // The grasp's background sync aligned the ref to a parked state event + // between this push's advertisement and its ref transaction. The ref + // is usually already where the push wants it, so a retry converges. + let reported = "remote: error: cannot lock ref 'refs/heads/main': is at \ + cac2ac91b6f5fb8dfcb6962785babc6e65350cb3 but expected \ + bc5e892aa84dc6240a5fbcd59367a4857d26f49b\n\ + To https://relay.ngit.dev/npub1owner/signed-test.git\n\ + ! [remote rejected] main -> main (incorrect old value provided)\n\ + error: failed to push some refs to 'https://relay.ngit.dev/npub1owner/signed-test.git'"; + assert!(is_transient_grasp_denial(reported)); + assert!(is_stale_advertisement_race(reported)); + + // Markers match independently of the surrounding git output. + assert!(is_stale_advertisement_race( + "cannot lock ref 'refs/heads/main'" + )); + assert!(is_stale_advertisement_race( + "! [remote rejected] main -> main (incorrect old value provided)" + )); + + // A purgatory denial is not a stale-advertisement race. + assert!(!is_stale_advertisement_race("No state events in purgatory")); + + // A real divergence is a different error and stays permanent. + assert!(!is_transient_grasp_denial( + " ! [rejected] main -> main (non-fast-forward)" + )); + } +} diff --git a/crates/signed_state/src/repo.rs b/crates/signed_state/src/repo.rs index 3d5e139..532ac6e 100644 --- a/crates/signed_state/src/repo.rs +++ b/crates/signed_state/src/repo.rs @@ -15,11 +15,10 @@ use signed_core::{ use signed_git::{Nip34Binding, PatchParser, Repo}; use signed_nostr::UniversalSigner; -use crate::backend::{ - Backend, BackendEvent, grasp_base_url, grasp06_prs_url, pr_clone_urls, require_relay_accepted, - user_grasp_list_servers, -}; +use crate::backend::{Backend, BackendEvent}; +use crate::bootstrap::user_grasp_list_servers; use crate::checkouts::CheckoutsStore; +use crate::push::{GraspPush, PushOutcome, grasp_base_url, grasp06_prs_url, pr_clone_urls}; use crate::repos::RepoListStore; /// Maximum size of one patch event. @@ -700,19 +699,12 @@ impl RepoStore { return; }; - let series: Vec = PatchParser::split_patch_series(&patch) - .into_iter() - .map(str::to_owned) - .collect(); + let series = PatchSeries::parse(&patch); - if let Some(oversized) = series - .iter() - .find(|part| part.len() > MAX_PATCH_EVENT_BYTES) - { + if let Some(oversized) = series.oversized_length() { self.last_error = Some(format!( "patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)", - oversized.len(), - MAX_PATCH_EVENT_BYTES + oversized, MAX_PATCH_EVENT_BYTES )); cx.notify(); return; @@ -720,11 +712,7 @@ impl RepoStore { // The tip of the series is its last commit. // `git format-patch` orders patches oldest first. - let Some(current_commit) = series - .last() - .and_then(|part| patch_current_commit(part)) - .and_then(|hex| hex.parse::().ok()) - else { + let Some(current_commit) = series.tip_commit() else { self.last_error = Some( "Patch must be `git format-patch` output with a `From ` header".into(), ); @@ -930,11 +918,10 @@ impl RepoStore { } } - let publish_result: Result = async { - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event) - } - .await; + let publish_result = { + let pusher = GraspPush::new(client.clone(), signer.clone()); + pusher.send_accepted(event).await + }; let pr_event = match publish_result { Ok(event) => event, @@ -1037,29 +1024,18 @@ impl RepoStore { return; } - let series: Vec = PatchParser::split_patch_series(&patch) - .into_iter() - .map(str::to_owned) - .collect(); - if let Some(oversized) = series - .iter() - .find(|part| part.len() > MAX_PATCH_EVENT_BYTES) - { + let series = PatchSeries::parse(&patch); + if let Some(oversized) = series.oversized_length() { self.last_error = Some(format!( "patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)", - oversized.len(), - MAX_PATCH_EVENT_BYTES + oversized, MAX_PATCH_EVENT_BYTES )); cx.notify(); return; } // The new tip of the PR is the last commit of the series. - let Some(current_commit) = series - .last() - .and_then(|part| patch_current_commit(part)) - .and_then(|hex| hex.parse::().ok()) - else { + let Some(current_commit) = series.tip_commit() else { self.last_error = Some( "Patch must be `git format-patch` output with a `From ` header".into(), ); @@ -1129,12 +1105,10 @@ impl RepoStore { } }; - let publish_result: Result = async { - let event = builder.finalize_async(&signer).await?; - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event) - } - .await; + let publish_result = { + let pusher = GraspPush::new(client.clone(), signer.clone()); + pusher.publish_one(builder).await + }; if let Err(e) = publish_result { return this.update(cx, |this, cx| { @@ -1216,38 +1190,10 @@ impl RepoStore { return self.action_error("Repository announcement is not loaded yet", cx); }; - self.pushing = true; - self.last_error = None; - self.last_push_warning = None; - cx.notify(); - let backend = Backend::global(cx); let push = backend.update(cx, |backend, cx| backend.push_repository(announcement, cx)); - cx.spawn(async move |this, cx| { - let result = push.await; - - this.update(cx, |this, cx| { - this.pushing = false; - - match &result { - Ok(outcome) => { - this.last_error = None; - // A push only some grasp servers accepted is a warning: - // the repo is out of sync on the rest until it is republished. - this.last_push_warning = outcome.partial_warning(); - } - Err(e) => { - this.last_error = Some(format!("Push failed: {e}")); - this.last_push_warning = None; - } - } - - cx.notify(); - })?; - - result.map(|_| ()) - }) + self.run_push(push, None, cx) } pub fn push_checkout( @@ -1273,17 +1219,30 @@ impl RepoStore { // The checkout may be on a side branch. let head = self.head.clone(); - self.pushing = true; - self.last_error = None; - self.last_push_warning = None; - cx.notify(); - - let checkouts = CheckoutsStore::global(cx); let backend = Backend::global(cx); let push = backend.update(cx, |backend, cx| { backend.push_checkout(announcement, path.clone(), head, cx) }); + self.run_push(push, Some((addr, path)), cx) + } + + /// Run a backend push task, tracking progress in [`Self::pushing`] and + /// the outcome in [`Self::last_error`] and [`Self::last_push_warning`]. + /// + /// `pushed_checkout` names the checkout whose ready-to-push statuses + /// should be recomputed after the remote moved. + fn run_push( + &mut self, + push: Task>, + pushed_checkout: Option<(RepoAddr, PathBuf)>, + cx: &mut Context, + ) -> Task> { + self.pushing = true; + self.last_error = None; + self.last_push_warning = None; + cx.notify(); + cx.spawn(async move |this, cx| { let result = push.await; @@ -1296,10 +1255,12 @@ impl RepoStore { // A push only some grasp servers accepted is a warning: // the repo is out of sync on the rest until it is republished. this.last_push_warning = outcome.partial_warning(); - // The remote moved, so recompute the ready-to-push statuses. - checkouts.update(cx, |store, cx| { - store.checkout_pushed(&addr, &path, cx); - }); + if let Some((addr, path)) = &pushed_checkout { + // The remote moved, so recompute the ready-to-push statuses. + CheckoutsStore::global(cx).update(cx, |store, cx| { + store.checkout_pushed(addr, path, cx); + }); + } } Err(e) => { this.last_error = Some(format!("Push failed: {e}")); @@ -1431,12 +1392,8 @@ impl RepoStore { }; let task: Task> = cx.spawn(async move |this, cx| { - let publish_result: Result = async { - let event = builder.finalize_async(&signer).await?; - let output = client.send_event(&event).broadcast().await?; - require_relay_accepted(output, event) - } - .await; + let pusher = GraspPush::new(client, signer); + let publish_result = pusher.publish_one(builder).await; if let Err(e) = publish_result { this.update(cx, |this, cx| { @@ -1497,6 +1454,46 @@ fn patch_current_commit(patch: &str) -> Option<&str> { hex.split_whitespace().next().filter(|hex| hex.len() == 40) } +/// A `git format-patch` series with the facts derived from its parts. +struct PatchSeries { + parts: Vec, +} + +impl PatchSeries { + fn parse(patch: &str) -> Self { + Self { + parts: PatchParser::split_patch_series(patch) + .into_iter() + .map(str::to_owned) + .collect(), + } + } + + /// The byte length of the first part over the NIP-34 size suggestion. + fn oversized_length(&self) -> Option { + self.parts + .iter() + .map(String::len) + .find(|length| *length > MAX_PATCH_EVENT_BYTES) + } + + /// The tip of the series is its last commit; + /// `git format-patch` orders patches oldest first. + fn tip_commit(&self) -> Option { + self.parts + .last() + .and_then(|part| patch_current_commit(part)) + .and_then(|hex| hex.parse::().ok()) + } + + fn commit_of(&self, index: usize) -> Option<&str> { + self.parts + .get(index) + .and_then(|part| patch_current_commit(part)) + .filter(|hex| hex.len() == 40) + } +} + /// Publish a `git format-patch` series as chained kind-1617 events. /// /// Returns the root event, the one a PR references. @@ -1507,15 +1504,16 @@ async fn publish_patch_series( addr: &RepoAddr, owner: PublicKey, euc: Option<&str>, - series: &[String], + series: &PatchSeries, first_marker: &str, reply_to: Option, ) -> Result { + let pusher = GraspPush::new(client.clone(), signer.clone()); let mut root: Option = None; let mut previous = reply_to; - for (ix, part) in series.iter().enumerate() { - let Some(commit) = patch_current_commit(part).filter(|hex| hex.len() == 40) else { + for (ix, part) in series.parts.iter().enumerate() { + let Some(commit) = series.commit_of(ix) else { return Err(anyhow::anyhow!( "patch {} of the series has no `From ` header", ix + 1 @@ -1557,9 +1555,7 @@ async fn publish_patch_series( } let builder = EventBuilder::new(Kind::GitPatch, part.clone()).tags(tags); - let event = builder.finalize_async(signer).await?; - let output = client.send_event(&event).broadcast().await?; - let event = require_relay_accepted(output, event)?; + let event = pusher.publish_one(builder).await?; if root.is_none() { root = Some(event.clone()); diff --git a/crates/workspace/src/views/inbox.rs b/crates/workspace/src/views/inbox.rs index 8f04a67..6761735 100644 --- a/crates/workspace/src/views/inbox.rs +++ b/crates/workspace/src/views/inbox.rs @@ -13,7 +13,7 @@ use gpui_component::{ActiveTheme, Icon, IconName, IconNamed, Sizable, StyledExt, use nostr::prelude::{Event, EventId, Kind, PublicKey, Timestamp}; use signed_core::{InboxItem, InboxReadState, RepoAddr}; use signed_state::{ - Backend, BackendEvent, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, query_inbox, + Backend, BackendEvent, Inbox, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, }; use signed_ui::{Avatar, CountBadge}; use utils::relative_time; @@ -199,7 +199,7 @@ impl InboxView { let state = self.state.clone(); - let work = cx.background_spawn(async move { query_inbox(&client, me, &state).await }); + let work = cx.background_spawn(async move { Inbox::query(&client, me, &state).await }); self.tasks.push(cx.spawn(async move |this, cx| { let (threads, unread_count) = match work.await { diff --git a/crates/workspace/src/views/pull_requests/detail.rs b/crates/workspace/src/views/pull_requests/detail.rs index 7b75eff..58b3b1c 100644 --- a/crates/workspace/src/views/pull_requests/detail.rs +++ b/crates/workspace/src/views/pull_requests/detail.rs @@ -23,7 +23,7 @@ use gpui_component::{ use nostr::prelude::{Event, EventId, Kind, Url}; use signed_core::{GitEvent, PullRequest, RepoAddr}; use signed_git::{FileCommit, PatchParser}; -use signed_state::{Backend, ProfileStore, RepoStore, ensure_repo_mirror}; +use signed_state::{Backend, Mirrors, ProfileStore, RepoStore}; use signed_ui::{Avatar, CountBadge, placeholder, status_badge}; use utils::{relative_time, relative_time_secs}; @@ -261,7 +261,7 @@ impl PullRequestDetailView { Some( cx.background_spawn(async move { - let repo = ensure_repo_mirror(&addr, &clone_urls)?; + let repo = Mirrors::ensure(&addr, &clone_urls)?; let gix_repo = repo.inner(); let workdir = gix_repo diff --git a/crates/workspace/src/views/pull_requests/new.rs b/crates/workspace/src/views/pull_requests/new.rs index 86e6744..d3cb0df 100644 --- a/crates/workspace/src/views/pull_requests/new.rs +++ b/crates/workspace/src/views/pull_requests/new.rs @@ -23,9 +23,7 @@ use gpui_component::{ use nostr::prelude::*; use signed_core::{Announcement, RepoAddr}; use signed_git::{GitCache, Repo}; -use signed_state::{ - Backend, CheckoutsStore, RepoListStore, RepoStore, ensure_repo_mirror, repo_mirror_path, -}; +use signed_state::{Backend, CheckoutsStore, Mirrors, RepoListStore, RepoStore}; use signed_ui::{CountBadge, placeholder, ref_selector_trigger}; use utils::middle_truncate; @@ -531,7 +529,7 @@ impl NewPullRequestView { let Some((base, _euc)) = self.base_repo(cx) else { return; }; - let mirror_path = repo_mirror_path(&base); + let mirror_path = Mirrors::path(&base); let namespace = GitCache::fork_namespace(&announcement); let clone_urls = announcement.clone.clone(); @@ -566,7 +564,7 @@ impl NewPullRequestView { let clone_urls = clone_urls.clone(); let mirror_path = mirror_path.clone(); async move { - ensure_repo_mirror(&base, &base_clone_urls)?; + Mirrors::ensure(&base, &base_clone_urls)?; let repo = Repo::open(&mirror_path)?; // Prune stale imports of any fork. Then import this fork's heads under its namespace. diff --git a/crates/workspace/src/views/repo/mod.rs b/crates/workspace/src/views/repo/mod.rs index 30a8dd4..a16c9f4 100644 --- a/crates/workspace/src/views/repo/mod.rs +++ b/crates/workspace/src/views/repo/mod.rs @@ -25,9 +25,8 @@ use nostr::prelude::{RelayUrl, ToBech32, Url}; use signed_core::{Announcement, RepoAddr, RepoStatus}; use signed_git::{FileCommit, GitCache, Repo}; use signed_state::{ - Backend, CheckoutStatus, CheckoutsStore, LocalReposStore, Nip34Binding, Nip34Kind, - ProfileStore, RepoListStore, RepoStore, ensure_repo_mirror, open_repo_mirror, - pr_proposes_checkout, + Backend, CheckoutStatus, CheckoutsStore, LocalReposStore, Mirrors, Nip34Binding, Nip34Kind, + ProfileStore, RepoListStore, RepoStore, }; use signed_ui::{ Avatar, CountBadge, DropdownButton, PixelAvatar, copy_row, menu_copy_row, ref_selector_trigger, @@ -359,7 +358,7 @@ impl RepoDetailView { let disk = { let addr = addr.clone(); cx.background_spawn(async move { - match open_repo_mirror(&addr)? { + match Mirrors::open(&addr)? { Some(repo) => Ok(Some(load_repo_data(&repo)?)), None => Ok(None), } @@ -376,7 +375,7 @@ impl RepoDetailView { let addr = addr.clone(); let clone_urls = clone_urls.clone(); cx.background_spawn(async move { - let repo = ensure_repo_mirror(&addr, &clone_urls)?; + let repo = Mirrors::ensure(&addr, &clone_urls)?; load_repo_data(&repo) }) .await @@ -403,7 +402,7 @@ impl RepoDetailView { let addr = addr.clone(); cx.background_spawn(async move { - let Some(repo) = open_repo_mirror(&addr)? else { + let Some(repo) = Mirrors::open(&addr)? else { return Ok::<_, Error>(None); }; @@ -1506,8 +1505,12 @@ impl RepoDetailView { continue; } for pr in &store.pull_requests { - if pr_proposes_checkout(pr, store.status_of(pr) == RepoStatus::Open, user, &status) - { + if CheckoutsStore::pr_proposes_checkout( + pr, + store.status_of(pr) == RepoStatus::Open, + user, + &status, + ) { continue 'status; } } diff --git a/crates/workspace/src/views/sidebar/mod.rs b/crates/workspace/src/views/sidebar/mod.rs index fd16d89..396f441 100644 --- a/crates/workspace/src/views/sidebar/mod.rs +++ b/crates/workspace/src/views/sidebar/mod.rs @@ -21,7 +21,7 @@ use nostr::prelude::RelayUrl; use signed_core::{Announcement, RepoAddr}; use signed_state::{ Backend, BackendEvent, CheckoutsStore, LocalReposStore, Nip34Binding, Nip34Kind, Profile, - ProfileStore, RepoListStore, ResolvedLocalRepo, resolve_local_repos, + ProfileStore, RepoListStore, ResolvedLocalRepo, }; use signed_ui::{Avatar, NavItem, PixelAvatar, title_bar_drag_handlers}; @@ -136,7 +136,7 @@ impl SidebarPanel { let local = LocalReposStore::global(cx).read(cx); let local_repos = - resolve_local_repos(&local.repos, &repo_list.announcements, &announcements); + LocalReposStore::resolve(&local.repos, &repo_list.announcements, &announcements); (announcements, local_repos, local.scanning) };