chore: refactor backend around domain types (#26)
Reviewed-on: #26
This commit was merged in pull request #26.
This commit is contained in:
+264
-1147
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,94 @@
|
||||
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;
|
||||
|
||||
pub const BOOTSTRAP_RELAYS: [&str; 2] = ["wss://relay.ditto.pub", "wss://index.ngit.dev"];
|
||||
pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"];
|
||||
|
||||
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<Filter>,
|
||||
) -> 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<Filter>> = 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<SyncSummary, Error> {
|
||||
ensure_bootstrap_relays(client).await?;
|
||||
|
||||
let output = client
|
||||
.sync(filter)
|
||||
.with(BOOTSTRAP_RELAYS)
|
||||
.opts(opts)
|
||||
.await?;
|
||||
|
||||
Ok(output.value)
|
||||
}
|
||||
|
||||
fn grasp_list_servers(event: &Event) -> Vec<RelayUrl> {
|
||||
event
|
||||
.tags
|
||||
.iter()
|
||||
.filter(|tag| tag.kind() == "g")
|
||||
.filter_map(|tag| tag.content())
|
||||
.filter_map(|url| RelayUrl::parse(url).ok())
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn latest_grasp_list_servers(events: Vec<Event>) -> Vec<RelayUrl> {
|
||||
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<Vec<RelayUrl>, Error> {
|
||||
let events: Vec<Event> = client
|
||||
.database()
|
||||
.query(Filters::grasp_list(user))
|
||||
.await?
|
||||
.into_iter()
|
||||
.collect();
|
||||
Ok(latest_grasp_list_servers(events))
|
||||
}
|
||||
@@ -7,20 +7,19 @@ use gpui::{App, AppContext, Context, Entity, Global, Subscription};
|
||||
use nostr::prelude::*;
|
||||
use settings::{CheckoutRecord, SettingsStore};
|
||||
use signed_core::{Announcement, RepoAddr};
|
||||
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;
|
||||
|
||||
const REFRESH_DEBOUNCE: Duration = Duration::from_millis(300);
|
||||
|
||||
/// How often the statuses are recomputed against the local refs.
|
||||
const LOCAL_POLL: Duration = Duration::from_secs(2);
|
||||
/// How often a full pass refreshes the remotes while any repository panel is open.
|
||||
const STATUS_POLL: Duration = Duration::from_secs(15);
|
||||
/// Remote refresh interval for the `ready to push` badges of the user's own repositories.
|
||||
const PUSH_POLL: Duration = Duration::from_secs(60);
|
||||
|
||||
const MAX_STATUS_CHECKOUTS: usize = 8;
|
||||
@@ -29,63 +28,37 @@ struct GlobalCheckoutsStore(Entity<CheckoutsStore>);
|
||||
|
||||
impl Global for GlobalCheckoutsStore {}
|
||||
|
||||
/// One associated local checkout of a repository.
|
||||
///
|
||||
/// Carries the git facts needed to suggest a pull request.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct CheckoutStatus {
|
||||
pub path: PathBuf,
|
||||
/// The branch checked out. A detached checkout is idle and yields no status.
|
||||
// A detached checkout is idle and yields no status.
|
||||
pub branch: String,
|
||||
/// Commit the branch points at, for tip-based PR dedupe.
|
||||
// For tip-based PR dedupe.
|
||||
pub head: String,
|
||||
/// What the branch is compared against.
|
||||
///
|
||||
/// It is `refs/remotes/origin/<branch>`, else `origin/HEAD` for new branches.
|
||||
// `refs/remotes/origin/<branch>`, else `origin/HEAD` for new branches.
|
||||
pub base: String,
|
||||
/// Commits in `base..branch`.
|
||||
///
|
||||
/// Zero-ahead checkouts are dropped, so this is always above zero.
|
||||
// Zero-ahead checkouts are dropped, so always above zero.
|
||||
pub ahead: u32,
|
||||
}
|
||||
|
||||
/// A remembered record, with the address already parsed.
|
||||
struct Remembered {
|
||||
path: PathBuf,
|
||||
addr: RepoAddr,
|
||||
last_used: u64,
|
||||
}
|
||||
|
||||
/// Global store of local-checkout associations and per-checkout statuses.
|
||||
pub struct CheckoutsStore {
|
||||
/// Checkout paths per announced repository.
|
||||
by_repo: HashMap<RepoAddr, Vec<PathBuf>>,
|
||||
/// Ready-to-contribute statuses of the requested repositories.
|
||||
///
|
||||
/// Those are the repository detail panels currently open.
|
||||
statuses: HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
/// Repositories whose statuses are recomputed on every input change.
|
||||
///
|
||||
/// Those are the repository detail panels currently open.
|
||||
status_requested: HashSet<RepoAddr>,
|
||||
/// Repositories whose `ready to push` statuses are recomputed on the same cycle.
|
||||
///
|
||||
/// The sidebar rows of the user's own repositories and their detail panels.
|
||||
push_requested: HashSet<RepoAddr>,
|
||||
/// The ready-to-push statuses of the requested own repositories.
|
||||
push_statuses: HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
/// Last announced head branch per requested repository.
|
||||
///
|
||||
/// A recompute defaults the base the same way.
|
||||
requested_head: HashMap<RepoAddr, Option<String>>,
|
||||
refresh: RefreshGate,
|
||||
/// True while the timer between a scheduled refresh and its run is pending.
|
||||
debounce_pending: bool,
|
||||
local_pending: bool,
|
||||
/// When the last full pass (with a remote refresh) completed.
|
||||
///
|
||||
/// The local pass runs a full pass again once this is older than the
|
||||
/// reconciliation cadence, so remote moves still land.
|
||||
// The local pass runs a full pass again once this is older than the
|
||||
// reconciliation cadence, so remote moves still land.
|
||||
last_full_sync: Option<Instant>,
|
||||
_subscriptions: Vec<Subscription>,
|
||||
}
|
||||
@@ -187,17 +160,10 @@ impl CheckoutsStore {
|
||||
});
|
||||
}
|
||||
|
||||
/// The associated checkouts of `addr`, freshest first.
|
||||
///
|
||||
/// Empty when none are known or the resolution has not run yet.
|
||||
pub fn associations_of(&self, addr: &RepoAddr) -> Vec<PathBuf> {
|
||||
self.by_repo.get(addr).cloned().unwrap_or_default()
|
||||
}
|
||||
|
||||
/// Ask for the `ready to contribute` statuses of `addr` to stay current.
|
||||
/// Called while the repository's detail panel is open.
|
||||
///
|
||||
/// `announced_head` is the announced HEAD branch, used to default the base.
|
||||
pub fn request_statuses(
|
||||
&mut self,
|
||||
addr: &RepoAddr,
|
||||
@@ -211,26 +177,18 @@ impl CheckoutsStore {
|
||||
self.refresh(cx);
|
||||
}
|
||||
|
||||
/// The ready-to-contribute statuses of `addr`.
|
||||
///
|
||||
/// Empty while none are known or nothing is ahead.
|
||||
pub fn ready_statuses_of(&self, addr: &RepoAddr) -> Vec<CheckoutStatus> {
|
||||
self.statuses.get(addr).cloned().unwrap_or_default()
|
||||
}
|
||||
|
||||
/// Ask for the `ready to push` statuses of `addr` to stay current.
|
||||
pub fn request_push_statuses(&mut self, addr: &RepoAddr, cx: &mut Context<Self>) {
|
||||
self.push_requested.insert(addr.clone());
|
||||
self.refresh(cx);
|
||||
}
|
||||
|
||||
/// The checkout at `path` was just pushed to the remote.
|
||||
///
|
||||
/// Its ready-to-push status is obsolete. Drop it from the cached statuses
|
||||
/// and notify observers right away, so the sidebar badge and the push
|
||||
/// banner update immediately instead of waiting for the next background
|
||||
/// pass, which re-scans and re-fetches the remote. The debounced refresh
|
||||
/// reconciles the remaining checkouts of the repository afterwards.
|
||||
// Drop the stale ready-to-push status and notify observers right away, so
|
||||
// the sidebar badge updates immediately instead of waiting for the next
|
||||
// background pass. The debounced refresh reconciles the remaining checkouts.
|
||||
pub fn checkout_pushed(&mut self, addr: &RepoAddr, path: &Path, cx: &mut Context<Self>) {
|
||||
let mut removed = false;
|
||||
|
||||
@@ -253,9 +211,6 @@ impl CheckoutsStore {
|
||||
self.request_push_statuses(addr, cx);
|
||||
}
|
||||
|
||||
/// The ready-to-push statuses of `addr`.
|
||||
///
|
||||
/// Empty while none are known or nothing is unpushed.
|
||||
pub fn push_statuses_of(&self, addr: &RepoAddr) -> Vec<CheckoutStatus> {
|
||||
self.push_statuses.get(addr).cloned().unwrap_or_default()
|
||||
}
|
||||
@@ -267,7 +222,6 @@ impl CheckoutsStore {
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
/// Re-resolve the associations and the requested statuses.
|
||||
pub fn refresh(&mut self, cx: &mut Context<Self>) {
|
||||
if self.debounce_pending || self.refresh.request() != RefreshRequest::Schedule {
|
||||
return;
|
||||
@@ -282,7 +236,6 @@ impl CheckoutsStore {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// One full resolve and apply cycle, the debounced entry point.
|
||||
fn run_refresh(&mut self, cx: &mut Context<Self>) {
|
||||
self.debounce_pending = false;
|
||||
self.refresh.begin();
|
||||
@@ -306,7 +259,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<String>)> = self
|
||||
.status_requested
|
||||
@@ -323,35 +276,38 @@ impl CheckoutsStore {
|
||||
let poll = !self.status_requested.is_empty() || !self.push_requested.is_empty();
|
||||
|
||||
let work = cx.background_spawn(async move {
|
||||
// Read the git facts of every scanned repository off the main thread.
|
||||
//
|
||||
// The facts are the origin URL and the root commit, both CLI reads.
|
||||
let mut facts: Vec<(PathBuf, Option<String>, Option<String>)> = Vec::new();
|
||||
for scanned in scanned.iter() {
|
||||
let path = &scanned.path;
|
||||
|
||||
// The browser's mirror clones share the announce URLs and EUCs. They are not user checkouts.
|
||||
// The browser's mirror clones are not user checkouts.
|
||||
if cache_root
|
||||
.as_ref()
|
||||
.is_some_and(|root| path.starts_with(root))
|
||||
{
|
||||
continue;
|
||||
}
|
||||
let origin = signed_git::origin_url(path).ok().flatten();
|
||||
let root = signed_git::root_commit(path).ok().flatten();
|
||||
let origin = Repo::open(path)
|
||||
.and_then(|repo| repo.origin_url())
|
||||
.ok()
|
||||
.flatten();
|
||||
let root = Repo::open(path)
|
||||
.and_then(|repo| repo.root_commit())
|
||||
.ok()
|
||||
.flatten();
|
||||
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<RepoAddr, Vec<PathBuf>> = associations
|
||||
.into_iter()
|
||||
.map(|(addr, paths)| (addr, paths.into_iter().filter(|p| p.is_dir()).collect()))
|
||||
.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))
|
||||
});
|
||||
@@ -360,7 +316,6 @@ impl CheckoutsStore {
|
||||
let (associations, statuses, push_statuses) = match work.await {
|
||||
Ok(results) => results,
|
||||
Err(_) => {
|
||||
// Git reads are best-effort, keep the last results.
|
||||
return this.update(cx, |this, cx| {
|
||||
this.refresh.abort();
|
||||
if poll {
|
||||
@@ -379,8 +334,7 @@ impl CheckoutsStore {
|
||||
this.statuses = statuses;
|
||||
this.push_statuses = push_statuses;
|
||||
|
||||
// Notify only when something actually changed, so observers
|
||||
// skip the no-op heartbeats.
|
||||
// Notify only when something actually changed.
|
||||
if associations_changed || statuses_changed || push_statuses_changed {
|
||||
cx.notify();
|
||||
}
|
||||
@@ -393,9 +347,6 @@ impl CheckoutsStore {
|
||||
this.update(cx, |this, cx| this.refresh(cx))?;
|
||||
}
|
||||
|
||||
// Restart the fast local pass so the freshly resolved
|
||||
// associations drive it. The pass itself decides when the next
|
||||
// full pass runs.
|
||||
this.update(cx, |this, cx| {
|
||||
if poll {
|
||||
this.schedule_local_pass(cx);
|
||||
@@ -407,7 +358,6 @@ impl CheckoutsStore {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// Schedule the fast local status pass, unless one is already pending.
|
||||
fn schedule_local_pass(&mut self, cx: &mut Context<Self>) {
|
||||
if self.local_pending {
|
||||
return;
|
||||
@@ -424,20 +374,17 @@ impl CheckoutsStore {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// The fast local status pass.
|
||||
fn local_tick(&mut self, cx: &mut Context<Self>) {
|
||||
// Nothing watched: the chain idles out until a new request restarts it.
|
||||
// Nothing watched: the pass idles until a new request restarts it.
|
||||
if self.status_requested.is_empty() && self.push_requested.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
// A full pass or a fresh request covers this tick, skip it.
|
||||
if self.refresh.running() || self.debounce_pending {
|
||||
self.schedule_local_pass(cx);
|
||||
return;
|
||||
}
|
||||
|
||||
// Open panels get the faster remote cadence.
|
||||
let cadence = if self.status_requested.is_empty() {
|
||||
PUSH_POLL
|
||||
} else {
|
||||
@@ -458,7 +405,6 @@ impl CheckoutsStore {
|
||||
self.schedule_local_pass(cx);
|
||||
}
|
||||
|
||||
/// Recompute the requested statuses against the tracking refs only.
|
||||
fn run_local_statuses(&mut self, cx: &mut Context<Self>) {
|
||||
let associations = self.by_repo.clone();
|
||||
|
||||
@@ -477,19 +423,18 @@ 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))
|
||||
});
|
||||
|
||||
let task: gpui::Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
|
||||
let Ok((statuses, push_statuses)) = work.await else {
|
||||
// Git reads are best-effort, keep the last results.
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
this.update(cx, |this, cx| {
|
||||
// A full pass or a fresh request will apply fresher data
|
||||
// (the tracking refs move only when a full pass fetches).
|
||||
// The tracking refs move only when a full pass fetches; a full
|
||||
// pass or a fresh request will apply fresher data.
|
||||
if this.refresh.running() || this.debounce_pending {
|
||||
return;
|
||||
}
|
||||
@@ -511,217 +456,201 @@ impl CheckoutsStore {
|
||||
}
|
||||
}
|
||||
|
||||
fn url_identity(url: &str) -> Option<(String, Option<u16>, String)> {
|
||||
let parsed = Url::parse(url).ok()?;
|
||||
let host = parsed.host_str()?.to_ascii_lowercase();
|
||||
let mut path = parsed.path().trim_matches('/').to_owned();
|
||||
if let Some(stripped) = path.strip_suffix(".git") {
|
||||
path = stripped.to_owned();
|
||||
}
|
||||
Some((host, parsed.port(), path))
|
||||
}
|
||||
impl CheckoutsStore {
|
||||
fn resolve_associations<'a>(
|
||||
remembered: &[Remembered],
|
||||
scanned: &[(PathBuf, Option<String>, Option<String>)],
|
||||
announcements: impl IntoIterator<Item = &'a Announcement>,
|
||||
) -> HashMap<RepoAddr, Vec<PathBuf>> {
|
||||
let announcements: Vec<&Announcement> = announcements.into_iter().collect();
|
||||
let mut out: HashMap<RepoAddr, Vec<PathBuf>> = HashMap::new();
|
||||
|
||||
fn same_repo_url(a: &str, b: &str) -> bool {
|
||||
match (url_identity(a), url_identity(b)) {
|
||||
(Some(a), Some(b)) => a == b,
|
||||
_ => a == b,
|
||||
}
|
||||
}
|
||||
|
||||
fn resolve_associations<'a>(
|
||||
remembered: &[Remembered],
|
||||
scanned: &[(PathBuf, Option<String>, Option<String>)],
|
||||
announcements: impl IntoIterator<Item = &'a Announcement>,
|
||||
) -> HashMap<RepoAddr, Vec<PathBuf>> {
|
||||
let announcements: Vec<&Announcement> = announcements.into_iter().collect();
|
||||
let mut out: HashMap<RepoAddr, Vec<PathBuf>> = 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<CheckoutStatus> {
|
||||
let repo = Repo::try_open(path)?;
|
||||
let branches = repo.branches().ok()?;
|
||||
|
||||
fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option<CheckoutStatus> {
|
||||
let branches = signed_git::worktree_branches(path).ok()?;
|
||||
if branches.is_empty() || repo.is_dirty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
if branches.is_empty() || signed_git::worktree_dirty(path) {
|
||||
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 = signed_git::worktree_current_branch(path)?;
|
||||
let head = signed_git::head_commit_id(path).ok().flatten()?;
|
||||
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())?;
|
||||
// Never fetches the checked-out refs; reads the tracking refs as-is.
|
||||
fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
|
||||
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 = signed_git::worktree_commits_ahead(path, &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<CheckoutStatus> {
|
||||
if signed_git::worktree_dirty(path) {
|
||||
return None;
|
||||
}
|
||||
let remote = format!("refs/remotes/origin/{branch}");
|
||||
|
||||
let branch = signed_git::worktree_current_branch(path)?;
|
||||
let head = signed_git::head_commit_id(path).ok().flatten()?;
|
||||
let origin = signed_git::origin_url(path).ok().flatten()?;
|
||||
|
||||
if fetch {
|
||||
signed_git::fetch_repo_refs(path, &[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 signed_git::worktree_ref_exists(path, &remote) {
|
||||
remote
|
||||
} else if signed_git::worktree_ref_exists(path, "refs/remotes/origin/HEAD") {
|
||||
"refs/remotes/origin/HEAD".to_owned()
|
||||
} else {
|
||||
return None;
|
||||
};
|
||||
|
||||
let ahead = signed_git::worktree_commits_ahead(path, &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<RepoAddr, Vec<PathBuf>>,
|
||||
requested: &[(RepoAddr, Option<String>)],
|
||||
push_requested: &[RepoAddr],
|
||||
fetch: bool,
|
||||
) -> (
|
||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
) {
|
||||
let mut statuses: HashMap<RepoAddr, Vec<CheckoutStatus>> = 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 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<CheckoutStatus> = 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,
|
||||
})
|
||||
}
|
||||
|
||||
fn compute_statuses(
|
||||
associations: &HashMap<RepoAddr, Vec<PathBuf>>,
|
||||
requested: &[(RepoAddr, Option<String>)],
|
||||
push_requested: &[RepoAddr],
|
||||
fetch: bool,
|
||||
) -> (
|
||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||
) {
|
||||
let mut statuses: HashMap<RepoAddr, Vec<CheckoutStatus>> = HashMap::new();
|
||||
for (addr, announced_head) in requested {
|
||||
let Some(paths) = associations.get(addr) else {
|
||||
continue;
|
||||
};
|
||||
|
||||
let list: Vec<CheckoutStatus> = 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<RepoAddr, Vec<CheckoutStatus>> = HashMap::new();
|
||||
for addr in push_requested {
|
||||
let Some(paths) = associations.get(addr) else {
|
||||
continue;
|
||||
};
|
||||
let mut push_statuses: HashMap<RepoAddr, Vec<CheckoutStatus>> = HashMap::new();
|
||||
for addr in push_requested {
|
||||
let Some(paths) = associations.get(addr) else {
|
||||
continue;
|
||||
};
|
||||
|
||||
let list: Vec<CheckoutStatus> = paths
|
||||
.iter()
|
||||
.take(MAX_STATUS_CHECKOUTS)
|
||||
.filter_map(|path| checkout_push_status(path, fetch))
|
||||
.collect();
|
||||
let list: Vec<CheckoutStatus> = 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)]
|
||||
mod tests {
|
||||
use std::process::Command;
|
||||
|
||||
use signed_core::{RepoAddr, repo_addr};
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn same_repo_url_ignores_the_transport_scheme() {
|
||||
// grasp announce vs https origin, with and without `.git`.
|
||||
assert!(same_repo_url(
|
||||
"grasp://relay.ngit.dev/npub1test/repo",
|
||||
"https://relay.ngit.dev/npub1test/repo.git"
|
||||
@@ -730,7 +659,6 @@ mod tests {
|
||||
"ws://localhost:8080/npub1test/repo",
|
||||
"http://localhost:8080/npub1test/repo"
|
||||
));
|
||||
// The port and the path matter.
|
||||
assert!(!same_repo_url(
|
||||
"wss://localhost:8081/npub1test/repo",
|
||||
"wss://localhost:8080/npub1test/repo"
|
||||
@@ -739,113 +667,15 @@ mod tests {
|
||||
"wss://host/npub1test/repo",
|
||||
"wss://host/npub1other/repo"
|
||||
));
|
||||
// Unparseable URLs compare literally.
|
||||
assert!(same_repo_url("/local/path", "/local/path"));
|
||||
assert!(!same_repo_url("/local/path", "/local/other"));
|
||||
}
|
||||
|
||||
fn remembered(path: &str, id: &str, last_used: u64) -> Remembered {
|
||||
Remembered {
|
||||
path: PathBuf::from(path),
|
||||
addr: addr(id),
|
||||
last_used,
|
||||
}
|
||||
}
|
||||
|
||||
fn scanned(
|
||||
path: &str,
|
||||
origin: Option<&str>,
|
||||
root: Option<&str>,
|
||||
) -> (PathBuf, Option<String>, Option<String>) {
|
||||
(
|
||||
PathBuf::from(path),
|
||||
origin.map(str::to_owned),
|
||||
root.map(str::to_owned),
|
||||
)
|
||||
}
|
||||
|
||||
const KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001";
|
||||
|
||||
fn owner() -> PublicKey {
|
||||
Keys::new(SecretKey::from_hex(KEY).expect("secret")).public_key()
|
||||
}
|
||||
|
||||
fn addr(id: &str) -> RepoAddr {
|
||||
repo_addr(owner(), id)
|
||||
}
|
||||
|
||||
/// Build one announcement by the fixed test owner.
|
||||
/// Takes `clone` URLs and an EUC.
|
||||
fn announcement(id: &str, clones: &[&str], euc: Option<&str>) -> Announcement {
|
||||
let keys = Keys::new(SecretKey::from_hex(KEY).expect("secret"));
|
||||
let mut tags = vec![Tag::parse(vec!["d", id]).expect("tag")];
|
||||
for url in clones {
|
||||
tags.push(Tag::parse(vec!["clone", *url]).expect("tag"));
|
||||
}
|
||||
if let Some(euc) = euc {
|
||||
tags.push(Tag::parse(vec!["r", euc, "euc"]).expect("tag"));
|
||||
}
|
||||
let event = EventBuilder::new(Kind::GitRepoAnnouncement, "")
|
||||
.tags(tags)
|
||||
.finalize(&keys)
|
||||
.expect("signed");
|
||||
Announcement::from_event(&event).expect("parsed")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn resolve_orders_remembered_freshest_first() {
|
||||
let announcements = vec![announcement("repo", &[], None)];
|
||||
let base = addr("repo");
|
||||
|
||||
let resolved = resolve_associations(
|
||||
&[
|
||||
remembered("/old", "repo", 100),
|
||||
remembered("/fresh", "repo", 200),
|
||||
remembered("/other", "unrelated", 300),
|
||||
],
|
||||
&[],
|
||||
&announcements,
|
||||
);
|
||||
|
||||
let paths = resolved.get(&base).expect("associations");
|
||||
assert_eq!(paths, &vec![PathBuf::from("/fresh"), PathBuf::from("/old")]);
|
||||
// Records for repositories without announcements stay inert.
|
||||
assert_eq!(resolved.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn resolve_deduplicates_paths_remembering_first() {
|
||||
let euc = "aa231c4c6a5777dc89b42207b499891a344add5c";
|
||||
let announcements = vec![announcement(
|
||||
"repo",
|
||||
&["https://host/npub1x/repo.git"],
|
||||
Some(euc),
|
||||
)];
|
||||
let base = addr("repo");
|
||||
|
||||
// The same path is both remembered and scanned, its origin matches.
|
||||
// The remembered occurrence wins and the path is listed once.
|
||||
let resolved = resolve_associations(
|
||||
&[remembered("/shared", "repo", 100)],
|
||||
&[
|
||||
scanned("/shared", Some("https://host/npub1x/repo"), None),
|
||||
scanned("/scanned-only", Some("https://host/npub1x/repo.git"), None),
|
||||
],
|
||||
&announcements,
|
||||
);
|
||||
|
||||
let paths = resolved.get(&base).expect("associations");
|
||||
assert_eq!(
|
||||
paths,
|
||||
&vec![PathBuf::from("/shared"), PathBuf::from("/scanned-only")]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checkout_status_reports_ahead_branches_only() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let path = dir.path().join("repo");
|
||||
let _initial = signed_git::init_repository(&path, "My Repo", "").expect("init");
|
||||
let _initial = Repo::init(&path, "My Repo", "").expect("init");
|
||||
let run = |args: &[&str]| {
|
||||
let status = Command::new("git")
|
||||
.current_dir(&path)
|
||||
@@ -864,35 +694,28 @@ mod tests {
|
||||
run(&["commit", "-m", message]);
|
||||
};
|
||||
|
||||
// A feature branch ahead of main, ready to contribute.
|
||||
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);
|
||||
assert_eq!(status.head.len(), 40);
|
||||
|
||||
// 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]
|
||||
fn checkout_push_status_counts_unpushed_commits_only() {
|
||||
// The `grasp remote` is a plain repository the checkout clones from.
|
||||
// Its origin URL is a local path, so the whole cycle runs offline.
|
||||
// Git refuses pushes to a checked-out branch by default.
|
||||
// Act like a grasp server and allow them.
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let remote = dir.path().join("remote");
|
||||
signed_git::init_repository(&remote, "My Repo", "").expect("init");
|
||||
Repo::init(&remote, "My Repo", "").expect("init");
|
||||
let config = Command::new("git")
|
||||
.args(["config", "receive.denyCurrentBranch", "ignore"])
|
||||
.current_dir(&remote)
|
||||
@@ -927,29 +750,23 @@ 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);
|
||||
assert_eq!(status.head.len(), 40);
|
||||
|
||||
// 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.
|
||||
let remote_run = |args: &[&str]| {
|
||||
let status = Command::new("git")
|
||||
.current_dir(&remote)
|
||||
@@ -965,6 +782,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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,40 +2,40 @@ use std::path::PathBuf;
|
||||
use std::sync::OnceLock;
|
||||
|
||||
use anyhow::Result;
|
||||
use gix::Repository;
|
||||
use nostr::prelude::Url;
|
||||
use signed_core::RepoAddr;
|
||||
use signed_git::GitCache;
|
||||
use signed_git::{GitCache, Repo};
|
||||
|
||||
static GIT_CACHE: OnceLock<GitCache> = OnceLock::new();
|
||||
|
||||
fn git_cache() -> &'static GitCache {
|
||||
GIT_CACHE
|
||||
.get()
|
||||
.expect("git cache is initialized by signed_state::init")
|
||||
}
|
||||
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 {
|
||||
pub fn install(root: impl Into<PathBuf>) {
|
||||
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<Option<Repository>> {
|
||||
git_cache().open(addr)
|
||||
}
|
||||
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<U: AsRef<str>>(addr: &RepoAddr, clone_urls: &[U]) -> Result<Repository> {
|
||||
git_cache().ensure_clone(addr, clone_urls)
|
||||
}
|
||||
pub fn path(addr: &RepoAddr) -> PathBuf {
|
||||
Self::cache().repo_path(addr)
|
||||
}
|
||||
|
||||
pub(crate) fn set_git_cache(root: impl Into<PathBuf>) {
|
||||
if GIT_CACHE.set(GitCache::new(root.into())).is_err() {
|
||||
log::warn!("git cache root is already set, keeping the first one");
|
||||
pub fn open(addr: &RepoAddr) -> Result<Option<Repo>> {
|
||||
Self::cache().open(addr)
|
||||
}
|
||||
|
||||
pub fn ensure(addr: &RepoAddr, clone_urls: &[Url]) -> Result<Repo> {
|
||||
Self::cache().ensure_clone(addr, clone_urls)
|
||||
}
|
||||
}
|
||||
|
||||
+131
-160
@@ -3,11 +3,10 @@ use std::collections::{HashMap, HashSet};
|
||||
use anyhow::Error;
|
||||
use gpui::{AppContext, Context, Task};
|
||||
use nostr_sdk::prelude::*;
|
||||
use signed_core::{Deletions, InboxItem, InboxReadState, filters, inbox};
|
||||
use signed_core::{Deletions, Filters, GitEvent, InboxItem, InboxReadState, inbox};
|
||||
|
||||
use crate::backend::Backend;
|
||||
|
||||
/// The user's persisted inbox read state.
|
||||
#[derive(Default)]
|
||||
pub struct Inbox {
|
||||
state: InboxReadState,
|
||||
@@ -15,7 +14,6 @@ pub struct Inbox {
|
||||
}
|
||||
|
||||
impl Inbox {
|
||||
/// The current read/archive cutoffs.
|
||||
pub fn state(&self) -> &InboxReadState {
|
||||
&self.state
|
||||
}
|
||||
@@ -24,41 +22,6 @@ impl Inbox {
|
||||
self.loaded
|
||||
}
|
||||
|
||||
pub fn mark_read(
|
||||
&mut self,
|
||||
group: &[Event],
|
||||
all: &[Event],
|
||||
me: PublicKey,
|
||||
cx: &mut Context<Self>,
|
||||
) {
|
||||
for event in group {
|
||||
self.state.mark_read(event);
|
||||
}
|
||||
self.state.advance_read(all, me, Timestamp::now());
|
||||
self.persist(cx);
|
||||
cx.notify();
|
||||
}
|
||||
|
||||
/// Archived events are always read too.
|
||||
pub fn mark_archived(
|
||||
&mut self,
|
||||
group: &[Event],
|
||||
all: &[Event],
|
||||
me: PublicKey,
|
||||
cx: &mut Context<Self>,
|
||||
) {
|
||||
for event in group {
|
||||
self.state.mark_archived(event);
|
||||
self.state.mark_read(event);
|
||||
}
|
||||
|
||||
let now = Timestamp::now();
|
||||
self.state.advance_archived(all, me, now);
|
||||
self.state.advance_read(all, me, now);
|
||||
self.persist(cx);
|
||||
cx.notify();
|
||||
}
|
||||
|
||||
pub fn mark_all_read(&mut self, all: &[Event], me: PublicKey, cx: &mut Context<Self>) {
|
||||
self.state.mark_all_read(all, me, Timestamp::now());
|
||||
self.persist(cx);
|
||||
@@ -71,7 +34,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;
|
||||
@@ -102,7 +65,6 @@ impl Inbox {
|
||||
cx.notify();
|
||||
}
|
||||
|
||||
/// Sign the state with a random key and store it locally.
|
||||
fn persist(&mut self, cx: &mut Context<Self>) {
|
||||
let backend = Backend::global(cx);
|
||||
let (me, client) = {
|
||||
@@ -117,7 +79,7 @@ impl Inbox {
|
||||
let state = self.state.clone();
|
||||
|
||||
let task: Task<Result<(), Error>> = 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(())
|
||||
@@ -125,135 +87,144 @@ impl Inbox {
|
||||
|
||||
task.detach();
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn query_inbox(
|
||||
client: &Client,
|
||||
me: PublicKey,
|
||||
state: &InboxReadState,
|
||||
) -> Result<(Vec<InboxItem>, usize), Error> {
|
||||
let deletion_events = client.database().query(filters::deletions()).await?;
|
||||
let deletions = Deletions::from_events(deletion_events);
|
||||
pub async fn query(
|
||||
client: &Client,
|
||||
me: PublicKey,
|
||||
state: &InboxReadState,
|
||||
) -> Result<(Vec<InboxItem>, 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) || !filters::is_git_activity(&event) {
|
||||
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<Option<InboxReadState>, 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<Event>, HashMap<EventId, Event>), Error> {
|
||||
let mut notifications: Vec<Event> = Vec::new();
|
||||
let mut by_id: HashMap<EventId, Event> = 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<EventId> = notifications.iter().flat_map(event_references).collect();
|
||||
let mut seen: HashSet<EventId> = 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 {
|
||||
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<Option<InboxReadState>, 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<Item = EventId> + '_ {
|
||||
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())
|
||||
})
|
||||
}
|
||||
|
||||
// Signed with a random key: the state is local-only, authorship does not
|
||||
// matter and no key material needs to be kept.
|
||||
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(())
|
||||
}
|
||||
|
||||
async fn fetch_notifications(
|
||||
client: &Client,
|
||||
me: PublicKey,
|
||||
deletions: &Deletions,
|
||||
) -> Result<(Vec<Event>, HashMap<EventId, Event>), Error> {
|
||||
let mut notifications: Vec<Event> = Vec::new();
|
||||
let mut by_id: HashMap<EventId, Event> = 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<EventId> = notifications
|
||||
.iter()
|
||||
.flat_map(Self::event_references)
|
||||
.collect();
|
||||
|
||||
let mut seen: HashSet<EventId> = 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))
|
||||
}
|
||||
|
||||
fn event_references(event: &Event) -> impl Iterator<Item = EventId> + '_ {
|
||||
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())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,31 +1,33 @@
|
||||
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, local_repo_addr, 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::{RepoActivityCounts, RepoListStore};
|
||||
pub use signed_git::{GraspSignals, LocalRepo, Nip34Binding, Nip34Kind};
|
||||
use signed_nostr::new_backend;
|
||||
pub use repos::RepoListStore;
|
||||
pub use signed_git::{Nip34Binding, Nip34Kind};
|
||||
use signed_nostr::NostrBackend;
|
||||
|
||||
#[cfg(not(target_arch = "wasm32"))]
|
||||
pub fn init(
|
||||
db_path: impl AsRef<Path>,
|
||||
repos_root: impl Into<PathBuf>,
|
||||
@@ -35,41 +37,17 @@ pub fn init(
|
||||
// rustls uses the `aws_lc_rs` provider by default.
|
||||
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
|
||||
|
||||
// Initialize the nostr client and signer
|
||||
let (client, signer) = cx.foreground_executor().block_on(async move {
|
||||
let path = db_path.as_ref().to_path_buf();
|
||||
new_backend(path)
|
||||
.await
|
||||
.expect("failed to initialize nostr backend")
|
||||
});
|
||||
let backend = cx
|
||||
.foreground_executor()
|
||||
.block_on(async move { NostrBackend::open(db_path.as_ref()).await });
|
||||
let backend = backend.expect("failed to initialize nostr backend");
|
||||
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);
|
||||
|
||||
// Set global stores for the profile
|
||||
ProfileStore::set_global(cx.new(ProfileStore::new), cx);
|
||||
|
||||
// Set global stores for the repo list and local repos
|
||||
RepoListStore::set_global(cx.new(RepoListStore::new), cx);
|
||||
|
||||
// Set global stores for the local repos
|
||||
LocalReposStore::set_global(cx.new(|cx| LocalReposStore::new(scan_paths, cx)), cx);
|
||||
|
||||
// Set global stores for the checkouts
|
||||
CheckoutsStore::set_global(cx.new(CheckoutsStore::new), cx);
|
||||
}
|
||||
|
||||
/// Initialize the backend with an in-memory database on wasm.
|
||||
#[cfg(target_arch = "wasm32")]
|
||||
pub fn init(cx: &mut App) {
|
||||
let (client, signer) = new_backend().expect("failed to initialize nostr backend");
|
||||
set_git_cache(PathBuf::new());
|
||||
Backend::set_global(cx.new(|cx| Backend::new(client, signer, cx)), cx);
|
||||
ProfileStore::set_global(cx.new(ProfileStore::new), cx);
|
||||
RepoListStore::set_global(cx.new(RepoListStore::new), cx);
|
||||
LocalReposStore::set_global(cx.new(|cx| LocalReposStore::new(Vec::new(), cx)), cx);
|
||||
CheckoutsStore::set_global(cx.new(|cx| CheckoutsStore::new(cx)), cx);
|
||||
}
|
||||
|
||||
@@ -4,17 +4,15 @@ use std::sync::Arc;
|
||||
|
||||
use anyhow::Error;
|
||||
use gpui::{App, AppContext, Context, Entity, Global, SharedString, Task};
|
||||
use signed_core::{Announcement, RepoAddr, repo_addr};
|
||||
use signed_core::{Announcement, RepoAddr};
|
||||
use signed_git::{LocalRepo, Nip34Binding, find_git_repos};
|
||||
|
||||
struct GlobalLocalReposStore(Entity<LocalReposStore>);
|
||||
|
||||
impl Global for GlobalLocalReposStore {}
|
||||
|
||||
/// Store of the git repositories discovered under a set of scan paths.
|
||||
pub struct LocalReposStore {
|
||||
pub roots: Arc<Vec<PathBuf>>,
|
||||
/// Git repositories discovered under [`Self::roots`], sorted by path.
|
||||
pub repos: Arc<Vec<LocalRepo>>,
|
||||
pub scanning: bool,
|
||||
scan_dirty: bool,
|
||||
@@ -45,7 +43,6 @@ impl LocalReposStore {
|
||||
}
|
||||
}
|
||||
|
||||
/// Forget a repository that has just been published to NIP-34.
|
||||
pub fn remove(&mut self, path: &Path, cx: &mut Context<Self>) {
|
||||
self.repos = Arc::new(
|
||||
self.repos
|
||||
@@ -94,7 +91,6 @@ impl LocalReposStore {
|
||||
dirty
|
||||
})?;
|
||||
|
||||
// Scans requested while this one ran are coalesced into one follow-up scan.
|
||||
if again {
|
||||
this.update(cx, |this, cx| this.rescan(cx))?;
|
||||
}
|
||||
@@ -106,28 +102,22 @@ impl LocalReposStore {
|
||||
}
|
||||
}
|
||||
|
||||
/// The NIP-34 coordinate a repository's detection resolved, when both the owner
|
||||
/// and the identifier were recovered.
|
||||
pub fn local_repo_addr(repo: &LocalRepo) -> Option<RepoAddr> {
|
||||
let binding = repo.nip34.as_ref()?;
|
||||
let owner = binding.owner?;
|
||||
let identifier = binding.identifier.as_deref()?;
|
||||
|
||||
Some(repo_addr(owner, identifier))
|
||||
Some(RepoAddr::new(owner, identifier))
|
||||
}
|
||||
|
||||
/// A scanned repository resolved against the known announcements.
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ResolvedLocalRepo {
|
||||
pub path: PathBuf,
|
||||
/// `None` for a plain repository.
|
||||
pub nip34: Option<Nip34Binding>,
|
||||
/// The known announcement this repository is bound to, when one matched.
|
||||
pub announcement: Option<Announcement>,
|
||||
}
|
||||
|
||||
impl ResolvedLocalRepo {
|
||||
/// The repository's directory name, or `Untitled` when the path has none.
|
||||
pub fn name(&self) -> SharedString {
|
||||
self.path
|
||||
.file_name()
|
||||
@@ -136,41 +126,42 @@ impl ResolvedLocalRepo {
|
||||
}
|
||||
}
|
||||
|
||||
/// Resolve the scanned repositories against the known announcements.
|
||||
pub fn resolve_local_repos(
|
||||
repos: &[LocalRepo],
|
||||
known: &[Announcement],
|
||||
own: &[Announcement],
|
||||
) -> Vec<ResolvedLocalRepo> {
|
||||
let shown: HashSet<RepoAddr> = own.iter().map(Announcement::addr).collect();
|
||||
impl LocalReposStore {
|
||||
pub fn resolve(
|
||||
repos: &[LocalRepo],
|
||||
known: &[Announcement],
|
||||
own: &[Announcement],
|
||||
) -> Vec<ResolvedLocalRepo> {
|
||||
let shown: HashSet<RepoAddr> = 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 +213,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,61 +221,9 @@ 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));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_unmatched_repository_keeps_its_binding() {
|
||||
let repo = bound(KEY, "unlisted");
|
||||
|
||||
let resolved = resolve_local_repos(&[repo], &[], &[]);
|
||||
|
||||
assert_eq!(resolved.len(), 1);
|
||||
assert!(resolved[0].announcement.is_none());
|
||||
assert_eq!(
|
||||
resolved[0].nip34.as_ref().map(|binding| binding.kind),
|
||||
Some(Nip34Kind::Initialized)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_plain_repository_is_kept_without_a_binding() {
|
||||
let repo = LocalRepo {
|
||||
path: PathBuf::from("plain"),
|
||||
nip34: None,
|
||||
};
|
||||
|
||||
let resolved = resolve_local_repos(&[repo], &[], &[]);
|
||||
|
||||
assert_eq!(resolved.len(), 1);
|
||||
assert!(resolved[0].nip34.is_none());
|
||||
assert!(resolved[0].announcement.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_name_is_the_directory_name() {
|
||||
let repo = LocalRepo {
|
||||
path: PathBuf::from("/tmp/my-repo"),
|
||||
nip34: None,
|
||||
};
|
||||
|
||||
let resolved = resolve_local_repos(&[repo], &[], &[]);
|
||||
|
||||
assert_eq!(resolved[0].name(), SharedString::from("my-repo"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_path_without_a_directory_name_is_untitled() {
|
||||
let repo = LocalRepo {
|
||||
path: PathBuf::from("/"),
|
||||
nip34: None,
|
||||
};
|
||||
|
||||
let resolved = resolve_local_repos(&[repo], &[], &[]);
|
||||
|
||||
assert_eq!(resolved[0].name(), SharedString::from("Untitled"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,16 +10,13 @@ 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);
|
||||
/// Max authors per profile request, keeping each filter within relay limits.
|
||||
const REQUEST_CHUNK: usize = 100;
|
||||
/// Recent profiles prefetched at startup and read back from the cache.
|
||||
const WARM_LIMIT: usize = 500;
|
||||
|
||||
/// A user profile as plain data for the UI, from the kind-0 metadata.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Profile {
|
||||
public_key: PublicKey,
|
||||
@@ -42,7 +39,7 @@ impl Profile {
|
||||
&self.metadata
|
||||
}
|
||||
|
||||
/// Display name, falling back to `name`, then a shortened npub.
|
||||
// Falls back to `name`, then a shortened npub.
|
||||
pub fn name(&self) -> SharedString {
|
||||
if let Some(display_name) = self.metadata.display_name.as_ref()
|
||||
&& !display_name.is_empty()
|
||||
@@ -68,14 +65,9 @@ impl Profile {
|
||||
}
|
||||
}
|
||||
|
||||
/// Global profile cache.
|
||||
///
|
||||
/// Profiles are fetched in batches and kept as plain data.
|
||||
pub struct ProfileStore {
|
||||
profiles: HashMap<PublicKey, Profile>,
|
||||
/// Public keys requested this session, main thread only.
|
||||
seen: RefCell<HashSet<PublicKey>>,
|
||||
/// Sender for queuing fetch requests, batched by a background task.
|
||||
sender: Sender<PublicKey>,
|
||||
_subscription: Subscription,
|
||||
}
|
||||
@@ -125,9 +117,7 @@ impl ProfileStore {
|
||||
}
|
||||
}
|
||||
|
||||
/// Get a profile.
|
||||
///
|
||||
/// Returns a placeholder with default metadata. Queues a fetch when the profile is not cached yet.
|
||||
// Returns a placeholder until fetched; queues a fetch when uncached.
|
||||
pub fn get(&self, public_key: &PublicKey) -> Profile {
|
||||
if let Some(profile) = self.profiles.get(public_key) {
|
||||
return profile.clone();
|
||||
@@ -178,7 +168,6 @@ impl ProfileStore {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// Re-read the latest metadata of `authors` from the local database in one query.
|
||||
fn apply_authors(&mut self, authors: Vec<PublicKey>, cx: &mut Context<Self>) {
|
||||
if authors.is_empty() {
|
||||
return;
|
||||
@@ -231,13 +220,11 @@ impl ProfileStore {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// Re-read the latest metadata of every requested author from the local database.
|
||||
fn apply_seen(&mut self, cx: &mut Context<Self>) {
|
||||
let authors: Vec<PublicKey> = self.seen.borrow().iter().copied().collect();
|
||||
self.apply_authors(authors, cx);
|
||||
}
|
||||
|
||||
/// Fetch metadata for requested authors in batches, debounced to collect requests.
|
||||
async fn handle_requests(
|
||||
this: WeakEntity<ProfileStore>,
|
||||
client: &Client,
|
||||
|
||||
@@ -0,0 +1,532 @@
|
||||
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;
|
||||
|
||||
const GRASP_RETRY_DELAY: Duration = Duration::from_secs(1);
|
||||
|
||||
// `ws://` grasp servers, like ngit, use `http://<host>`; secure relays map to
|
||||
// `https://<host>`.
|
||||
pub(crate) fn grasp_base_url(relay: &RelayUrl) -> Option<String> {
|
||||
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<Url> {
|
||||
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")
|
||||
}
|
||||
|
||||
// The author's GRASP-06 `/prs/` URLs come first.
|
||||
pub(crate) fn pr_clone_urls(prs_urls: Vec<Url>, base_clone_urls: Vec<Url>) -> Vec<Url> {
|
||||
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,
|
||||
reason: Option<String>,
|
||||
}
|
||||
|
||||
impl GraspServer {
|
||||
fn ok(relay: RelayUrl) -> Self {
|
||||
Self {
|
||||
relay,
|
||||
reason: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn failed(relay: RelayUrl, reason: impl Into<String>) -> Self {
|
||||
Self {
|
||||
relay,
|
||||
reason: Some(reason.into()),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn relay(&self) -> &RelayUrl {
|
||||
&self.relay
|
||||
}
|
||||
|
||||
pub fn reason(&self) -> Option<&str> {
|
||||
self.reason.as_deref()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct PushOutcome {
|
||||
pub servers: Vec<GraspServer>,
|
||||
// The newest state event a grasp relay accepted for this push; broadcast
|
||||
// to the other relays once a git server holds the data.
|
||||
pub state_event: Option<Event>,
|
||||
}
|
||||
|
||||
impl PushOutcome {
|
||||
pub fn accepted(&self) -> usize {
|
||||
self.servers
|
||||
.iter()
|
||||
.filter(|server| server.reason.is_none())
|
||||
.count()
|
||||
}
|
||||
|
||||
fn failing(&self) -> impl Iterator<Item = &GraspServer> {
|
||||
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::<Vec<_>>()
|
||||
.join("; ")
|
||||
}
|
||||
|
||||
// `None` when every server accepted the push or nothing was pushed.
|
||||
pub fn partial_warning(&self) -> Option<String> {
|
||||
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")
|
||||
}
|
||||
|
||||
// All staged events carry the same refs; the newest timestamp wins on the relays.
|
||||
fn keep_newest(state_event: &mut Option<Event>, 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 }
|
||||
}
|
||||
|
||||
pub(crate) fn require_relay_accepted(
|
||||
output: SendEventOutput,
|
||||
event: Event,
|
||||
) -> Result<Event, Error> {
|
||||
if output.success.is_empty() && !output.failed.is_empty() {
|
||||
let reasons = output
|
||||
.failed
|
||||
.values()
|
||||
.cloned()
|
||||
.collect::<Vec<String>>()
|
||||
.join(", ");
|
||||
bail!("event not accepted by any relay: {reasons}");
|
||||
}
|
||||
|
||||
Ok(event)
|
||||
}
|
||||
|
||||
// A relay hiccup should not block sign-up, so failures are only logged.
|
||||
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}");
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn publish_one(&self, builder: EventBuilder) -> Result<Event, Error> {
|
||||
let event = builder.finalize_async(&self.signer).await?;
|
||||
self.send_accepted(event).await
|
||||
}
|
||||
|
||||
pub(crate) async fn send_accepted(&self, event: Event) -> Result<Event, Error> {
|
||||
let output = self.client.send_event(&event).broadcast().await?;
|
||||
Self::require_relay_accepted(output, 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(())
|
||||
}
|
||||
|
||||
// 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))
|
||||
}
|
||||
|
||||
// `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;
|
||||
// For the convergence probe 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 {
|
||||
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;
|
||||
|
||||
// A failed stage means the grasp never parked the state, so
|
||||
// the git push would be denied anyway: skip it.
|
||||
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() {
|
||||
} 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. 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"
|
||||
);
|
||||
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<_>>(),
|
||||
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() {
|
||||
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));
|
||||
|
||||
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"
|
||||
));
|
||||
|
||||
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() {
|
||||
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));
|
||||
|
||||
assert!(is_stale_advertisement_race(
|
||||
"cannot lock ref 'refs/heads/main'"
|
||||
));
|
||||
assert!(is_stale_advertisement_race(
|
||||
"! [remote rejected] main -> main (incorrect old value provided)"
|
||||
));
|
||||
|
||||
assert!(!is_stale_advertisement_race("No state events in purgatory"));
|
||||
|
||||
assert!(!is_transient_grasp_denial(
|
||||
" ! [rejected] main -> main (non-fast-forward)"
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,3 @@
|
||||
/// Refresh coalescing shared by the event stores.
|
||||
#[derive(Debug, Default)]
|
||||
pub struct RefreshGate {
|
||||
running: bool,
|
||||
@@ -7,9 +6,7 @@ pub struct RefreshGate {
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum RefreshRequest {
|
||||
/// No run covers the request, start one now.
|
||||
Schedule,
|
||||
/// A run is in flight and covers the request, fold it into a follow-up.
|
||||
Fold,
|
||||
}
|
||||
|
||||
@@ -31,59 +28,13 @@ impl RefreshGate {
|
||||
self.running = true;
|
||||
}
|
||||
|
||||
/// The run ended. Whether a request arrived while it ran.
|
||||
pub fn finish(&mut self) -> bool {
|
||||
self.running = false;
|
||||
std::mem::take(&mut self.dirty)
|
||||
}
|
||||
|
||||
/// The run was abandoned, e.g. on error. Pending follow-up requests survive.
|
||||
// Pending follow-up requests survive an abandoned run.
|
||||
pub fn abort(&mut self) {
|
||||
self.running = false;
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn a_request_while_running_folds_into_a_follow_up() {
|
||||
let mut gate = RefreshGate::default();
|
||||
gate.begin();
|
||||
|
||||
assert_eq!(gate.request(), RefreshRequest::Fold);
|
||||
assert!(gate.finish());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_request_without_a_run_schedules() {
|
||||
let mut gate = RefreshGate::default();
|
||||
|
||||
assert_eq!(gate.request(), RefreshRequest::Schedule);
|
||||
assert!(!gate.running());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_request_after_a_run_schedules_again() {
|
||||
let mut gate = RefreshGate::default();
|
||||
gate.begin();
|
||||
assert_eq!(gate.request(), RefreshRequest::Fold);
|
||||
assert!(gate.finish());
|
||||
|
||||
assert_eq!(gate.request(), RefreshRequest::Schedule);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn abort_keeps_the_pending_request() {
|
||||
let mut gate = RefreshGate::default();
|
||||
gate.begin();
|
||||
assert_eq!(gate.request(), RefreshRequest::Fold);
|
||||
|
||||
gate.abort();
|
||||
assert!(!gate.running());
|
||||
|
||||
gate.begin();
|
||||
assert!(gate.finish());
|
||||
}
|
||||
}
|
||||
|
||||
+177
-444
File diff suppressed because it is too large
Load Diff
@@ -5,7 +5,7 @@ use std::time::Duration;
|
||||
use anyhow::Error;
|
||||
use gpui::{App, AppContext, Context, Entity, Global, Subscription};
|
||||
use nostr_sdk::prelude::*;
|
||||
use signed_core::{Announcement, Deletions, RepoAddr, filters, repo_addr};
|
||||
use signed_core::{Announcement, Deletions, Filters, RepoAddr, filters};
|
||||
|
||||
use crate::backend::{Backend, BackendEvent};
|
||||
use crate::refresh::{RefreshGate, RefreshRequest};
|
||||
@@ -17,37 +17,27 @@ struct GlobalRepoListStore(Entity<RepoListStore>);
|
||||
|
||||
impl Global for GlobalRepoListStore {}
|
||||
|
||||
/// NIP-34 activity event counts per repository, ranking the explore list by popularity.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct RepoActivityCounts {
|
||||
/// Root `30611` issue events addressed to the repository.
|
||||
pub issues: u32,
|
||||
/// Root `3063` pull request events addressed to the repository.
|
||||
///
|
||||
/// PR updates are not new PRs and do not count.
|
||||
// PR updates are not new PRs and do not count.
|
||||
pub pull_requests: u32,
|
||||
/// `1617` patch events addressed to the repository.
|
||||
pub commits: u32,
|
||||
}
|
||||
|
||||
impl RepoActivityCounts {
|
||||
/// Total issues, pull requests and commits, the popularity ranking key.
|
||||
// The popularity ranking key.
|
||||
pub fn score(self) -> u32 {
|
||||
self.issues + self.pull_requests + self.commits
|
||||
}
|
||||
}
|
||||
|
||||
/// Store listing the discovered repository announcements, newest first.
|
||||
pub struct RepoListStore {
|
||||
/// Shared so views can clone the list per frame without a deep copy.
|
||||
// Shared so views can clone the list per frame without a deep copy.
|
||||
pub announcements: Arc<Vec<Announcement>>,
|
||||
/// Latest known activity timestamp per repository.
|
||||
pub last_activity: Arc<HashMap<RepoAddr, Timestamp>>,
|
||||
/// Issues, pull requests and commits per repository.
|
||||
///
|
||||
/// Used for the Popular ranking of the explore list.
|
||||
// For the Popular ranking of the explore list.
|
||||
pub counts: Arc<HashMap<RepoAddr, RepoActivityCounts>>,
|
||||
/// Own repositories whose state events were fetched from their announced relays.
|
||||
state_synced_repos: HashSet<RepoAddr>,
|
||||
refresh: RefreshGate,
|
||||
_subscription: Subscription,
|
||||
@@ -100,7 +90,6 @@ impl RepoListStore {
|
||||
}
|
||||
}
|
||||
|
||||
/// The announcements of `user`, newest first.
|
||||
pub fn announcements_of(&self, user: &PublicKey) -> Vec<Announcement> {
|
||||
self.announcements
|
||||
.iter()
|
||||
@@ -115,17 +104,16 @@ impl RepoListStore {
|
||||
backend.update(cx, |backend, cx| {
|
||||
backend.sync_bootstraps(
|
||||
vec![
|
||||
filters::all_announcements(),
|
||||
filters::all_states(),
|
||||
Filters::all_announcements(),
|
||||
Filters::all_states(),
|
||||
// Deletion requests, NIP-09/62, must be known before any announcement is shown.
|
||||
filters::deletions(),
|
||||
Filters::deletions(),
|
||||
],
|
||||
cx,
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
/// Fetch the state events of the user's own repositories.
|
||||
fn sync_own_repo_states(&mut self, cx: &mut Context<Self>) {
|
||||
let backend = Backend::global(cx);
|
||||
let Some(me) = backend.read(cx).current_user() else {
|
||||
@@ -144,15 +132,13 @@ impl RepoListStore {
|
||||
self.state_synced_repos.insert(addr.clone());
|
||||
|
||||
backend.update(cx, |backend, cx| {
|
||||
backend.connect_repo_relays(relays, vec![filters::state(&addr)], cx);
|
||||
backend.connect_repo_relays(relays, vec![addr.state_filter()], cx);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/// Re-query the local database.
|
||||
///
|
||||
/// Runs immediately. The backend pump already batches the relay events that
|
||||
/// trigger a refresh, so no per-store debounce is needed.
|
||||
// Runs immediately: the backend pump already batches the relay events
|
||||
// that trigger a refresh, so no per-store debounce is needed.
|
||||
pub fn refresh(&mut self, cx: &mut Context<Self>) {
|
||||
if self.refresh.request() != RefreshRequest::Schedule {
|
||||
return;
|
||||
@@ -168,14 +154,14 @@ impl RepoListStore {
|
||||
let client = backend.read(cx).client();
|
||||
|
||||
let work = cx.background_spawn(async move {
|
||||
let filter = filters::all_announcements();
|
||||
let filter = Filters::all_announcements();
|
||||
let events = client.database().query(filter).await?;
|
||||
|
||||
let deletion_events = client.database().query(filters::deletions()).await?;
|
||||
let deletion_events = client.database().query(Filters::deletions()).await?;
|
||||
let deletions = Deletions::from_events(deletion_events);
|
||||
|
||||
// Dedup and sort off the main thread.
|
||||
// Only the final list crosses back into the entity.
|
||||
// Dedup and sort off the main thread; only the final list crosses
|
||||
// back into the entity.
|
||||
let mut by_repo: HashMap<RepoAddr, Announcement> = HashMap::new();
|
||||
|
||||
for event in events {
|
||||
@@ -213,14 +199,13 @@ impl RepoListStore {
|
||||
let Some(id) = event.tags.identifier() else {
|
||||
continue;
|
||||
};
|
||||
let addr = repo_addr(event.pubkey, id);
|
||||
let addr = RepoAddr::new(event.pubkey, id);
|
||||
let Some(entry) = last_activity.get_mut(&addr) else {
|
||||
continue;
|
||||
};
|
||||
*entry = (*entry).max(event.created_at);
|
||||
}
|
||||
|
||||
// Bound the activity query to a recent window.
|
||||
// Older repos fall back to their announcement or state timestamps.
|
||||
let activity_filter = Filter::new()
|
||||
.kinds(filters::ACTIVITY_KINDS)
|
||||
@@ -230,10 +215,11 @@ impl RepoListStore {
|
||||
if deletions.is_deleted(&event) {
|
||||
continue;
|
||||
}
|
||||
for addr in event.tags.coordinates() {
|
||||
if addr.kind != Kind::GitRepoAnnouncement {
|
||||
for coordinate in event.tags.coordinates() {
|
||||
if coordinate.kind != Kind::GitRepoAnnouncement {
|
||||
continue;
|
||||
}
|
||||
let addr = RepoAddr::from(coordinate.clone());
|
||||
let Some(entry) = last_activity.get_mut(&addr) else {
|
||||
continue;
|
||||
};
|
||||
@@ -241,8 +227,8 @@ impl RepoListStore {
|
||||
}
|
||||
}
|
||||
|
||||
// Popularity counts per repository, issues, pull requests and patches.
|
||||
// Unbounded, unlike the windowed activity query above, so totals are exact.
|
||||
// Unbounded, unlike the windowed activity query above, so totals
|
||||
// are exact.
|
||||
let mut counts: HashMap<RepoAddr, RepoActivityCounts> = HashMap::new();
|
||||
let count_filter =
|
||||
Filter::new().kinds([Kind::GitIssue, Kind::GitPullRequest, Kind::GitPatch]);
|
||||
@@ -251,12 +237,13 @@ impl RepoListStore {
|
||||
if deletions.is_deleted(&event) {
|
||||
continue;
|
||||
}
|
||||
for addr in event.tags.coordinates() {
|
||||
if addr.kind != Kind::GitRepoAnnouncement || !last_activity.contains_key(&addr)
|
||||
for coordinate in event.tags.coordinates() {
|
||||
if coordinate.kind != Kind::GitRepoAnnouncement
|
||||
|| !last_activity.contains_key(&RepoAddr::from(coordinate.clone()))
|
||||
{
|
||||
continue;
|
||||
}
|
||||
let entry = counts.entry(addr).or_default();
|
||||
let entry = counts.entry(RepoAddr::from(coordinate)).or_default();
|
||||
match event.kind {
|
||||
Kind::GitIssue => entry.issues += 1,
|
||||
Kind::GitPullRequest => entry.pull_requests += 1,
|
||||
@@ -272,7 +259,6 @@ impl RepoListStore {
|
||||
cx.spawn(async move |this, cx| {
|
||||
let (announcements, last_activity, counts) = match work.await {
|
||||
Ok(results) => results,
|
||||
// Database errors are transient, keep the last list.
|
||||
Err(_) => {
|
||||
return this.update(cx, |this, _cx| {
|
||||
this.refresh.abort();
|
||||
@@ -290,8 +276,6 @@ impl RepoListStore {
|
||||
this.refresh.finish()
|
||||
})?;
|
||||
|
||||
// Requests that arrived while the refresh was running.
|
||||
// They are coalesced into one follow-up refresh.
|
||||
if again {
|
||||
this.update(cx, |this, cx| this.refresh(cx))?;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user