chore: refactor backend around domain types #26
+189
-840
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,99 @@
|
|||||||
|
use std::collections::HashMap;
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
use anyhow::Error;
|
||||||
|
use nostr_connect::prelude::*;
|
||||||
|
use nostr_sdk::client::SyncSummary;
|
||||||
|
use nostr_sdk::prelude::*;
|
||||||
|
use signed_core::Filters;
|
||||||
|
|
||||||
|
/// Relays connected at startup, before any user-specific relay config is known.
|
||||||
|
pub const BOOTSTRAP_RELAYS: [&str; 2] = ["wss://relay.ditto.pub", "wss://index.ngit.dev"];
|
||||||
|
/// Relays used to index the user's NIP-65 relay list.
|
||||||
|
pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"];
|
||||||
|
|
||||||
|
/// Add and connect the startup relays.
|
||||||
|
async fn ensure_bootstrap_relays(client: &Client) -> Result<(), Error> {
|
||||||
|
for url in BOOTSTRAP_RELAYS {
|
||||||
|
client.add_relay(url).and_connect().await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
for url in INDEXER_RELAYS {
|
||||||
|
client
|
||||||
|
.add_relay(url)
|
||||||
|
.capabilities(RelayCapabilities::DISCOVERY)
|
||||||
|
.and_connect()
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn subscribe_bootstrap_only(
|
||||||
|
client: &Client,
|
||||||
|
filters: Vec<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)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The `g` tag servers of one kind-10317 grasp list event, in tag order.
|
||||||
|
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()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Grasp servers of the newest kind-10317 grasp list among `events`.
|
||||||
|
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))
|
||||||
|
}
|
||||||
@@ -11,7 +11,7 @@ use signed_git::Repo;
|
|||||||
use utils::same_repo_url;
|
use utils::same_repo_url;
|
||||||
|
|
||||||
use crate::backend::{Backend, BackendEvent};
|
use crate::backend::{Backend, BackendEvent};
|
||||||
use crate::git_store::repo_mirror_root;
|
use crate::git_store::Mirrors;
|
||||||
use crate::local_repos::LocalReposStore;
|
use crate::local_repos::LocalReposStore;
|
||||||
use crate::refresh::{RefreshGate, RefreshRequest};
|
use crate::refresh::{RefreshGate, RefreshRequest};
|
||||||
use crate::repos::RepoListStore;
|
use crate::repos::RepoListStore;
|
||||||
@@ -308,7 +308,7 @@ impl CheckoutsStore {
|
|||||||
|
|
||||||
let announcements = RepoListStore::global(cx).read(cx).announcements.clone();
|
let announcements = RepoListStore::global(cx).read(cx).announcements.clone();
|
||||||
let scanned = LocalReposStore::global(cx).read(cx).repos.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
|
let requested: Vec<(RepoAddr, Option<String>)> = self
|
||||||
.status_requested
|
.status_requested
|
||||||
@@ -350,7 +350,8 @@ impl CheckoutsStore {
|
|||||||
facts.push((path.clone(), origin, root));
|
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.
|
// Missing directories are stale records, drop them.
|
||||||
let associations: HashMap<RepoAddr, Vec<PathBuf>> = associations
|
let associations: HashMap<RepoAddr, Vec<PathBuf>> = associations
|
||||||
@@ -359,7 +360,7 @@ impl CheckoutsStore {
|
|||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
let (statuses, push_statuses) =
|
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))
|
Ok::<_, Error>((associations, statuses, push_statuses))
|
||||||
});
|
});
|
||||||
@@ -485,7 +486,7 @@ impl CheckoutsStore {
|
|||||||
|
|
||||||
let work = cx.background_spawn(async move {
|
let work = cx.background_spawn(async move {
|
||||||
let (statuses, push_statuses) =
|
let (statuses, push_statuses) =
|
||||||
compute_statuses(&associations, &requested, &push_requested, false);
|
CheckoutsStore::compute_statuses(&associations, &requested, &push_requested, false);
|
||||||
Ok::<_, Error>((statuses, push_statuses))
|
Ok::<_, Error>((statuses, push_statuses))
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -519,11 +520,12 @@ impl CheckoutsStore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn resolve_associations<'a>(
|
impl CheckoutsStore {
|
||||||
|
fn resolve_associations<'a>(
|
||||||
remembered: &[Remembered],
|
remembered: &[Remembered],
|
||||||
scanned: &[(PathBuf, Option<String>, Option<String>)],
|
scanned: &[(PathBuf, Option<String>, Option<String>)],
|
||||||
announcements: impl IntoIterator<Item = &'a Announcement>,
|
announcements: impl IntoIterator<Item = &'a Announcement>,
|
||||||
) -> HashMap<RepoAddr, Vec<PathBuf>> {
|
) -> HashMap<RepoAddr, Vec<PathBuf>> {
|
||||||
let announcements: Vec<&Announcement> = announcements.into_iter().collect();
|
let announcements: Vec<&Announcement> = announcements.into_iter().collect();
|
||||||
let mut out: HashMap<RepoAddr, Vec<PathBuf>> = HashMap::new();
|
let mut out: HashMap<RepoAddr, Vec<PathBuf>> = HashMap::new();
|
||||||
|
|
||||||
@@ -559,9 +561,9 @@ fn resolve_associations<'a>(
|
|||||||
}
|
}
|
||||||
|
|
||||||
out
|
out
|
||||||
}
|
}
|
||||||
|
|
||||||
fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option<CheckoutStatus> {
|
fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option<CheckoutStatus> {
|
||||||
let repo = Repo::try_open(path)?;
|
let repo = Repo::try_open(path)?;
|
||||||
let branches = repo.branches().ok()?;
|
let branches = repo.branches().ok()?;
|
||||||
|
|
||||||
@@ -589,10 +591,10 @@ fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option<Checkout
|
|||||||
base,
|
base,
|
||||||
ahead,
|
ahead,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The `ready to push` status of one checkout of the user's own repository.
|
/// The `ready to push` status of one checkout of the user's own repository.
|
||||||
fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
|
fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
|
||||||
let repo = Repo::try_open(path)?;
|
let repo = Repo::try_open(path)?;
|
||||||
if repo.is_dirty() {
|
if repo.is_dirty() {
|
||||||
return None;
|
return None;
|
||||||
@@ -628,18 +630,18 @@ fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
|
|||||||
base,
|
base,
|
||||||
ahead,
|
ahead,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Compute the requested statuses against the checkout paths of `associations`.
|
/// Compute the requested statuses against the checkout paths of `associations`.
|
||||||
fn compute_statuses(
|
fn compute_statuses(
|
||||||
associations: &HashMap<RepoAddr, Vec<PathBuf>>,
|
associations: &HashMap<RepoAddr, Vec<PathBuf>>,
|
||||||
requested: &[(RepoAddr, Option<String>)],
|
requested: &[(RepoAddr, Option<String>)],
|
||||||
push_requested: &[RepoAddr],
|
push_requested: &[RepoAddr],
|
||||||
fetch: bool,
|
fetch: bool,
|
||||||
) -> (
|
) -> (
|
||||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||||
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
HashMap<RepoAddr, Vec<CheckoutStatus>>,
|
||||||
) {
|
) {
|
||||||
let mut statuses: HashMap<RepoAddr, Vec<CheckoutStatus>> = HashMap::new();
|
let mut statuses: HashMap<RepoAddr, Vec<CheckoutStatus>> = HashMap::new();
|
||||||
for (addr, announced_head) in requested {
|
for (addr, announced_head) in requested {
|
||||||
let Some(paths) = associations.get(addr) else {
|
let Some(paths) = associations.get(addr) else {
|
||||||
@@ -649,7 +651,7 @@ fn compute_statuses(
|
|||||||
let list: Vec<CheckoutStatus> = paths
|
let list: Vec<CheckoutStatus> = paths
|
||||||
.iter()
|
.iter()
|
||||||
.take(MAX_STATUS_CHECKOUTS)
|
.take(MAX_STATUS_CHECKOUTS)
|
||||||
.filter_map(|path| checkout_status(path, announced_head.as_deref()))
|
.filter_map(|path| Self::checkout_status(path, announced_head.as_deref()))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
if !list.is_empty() {
|
if !list.is_empty() {
|
||||||
@@ -666,7 +668,7 @@ fn compute_statuses(
|
|||||||
let list: Vec<CheckoutStatus> = paths
|
let list: Vec<CheckoutStatus> = paths
|
||||||
.iter()
|
.iter()
|
||||||
.take(MAX_STATUS_CHECKOUTS)
|
.take(MAX_STATUS_CHECKOUTS)
|
||||||
.filter_map(|path| checkout_push_status(path, fetch))
|
.filter_map(|path| Self::checkout_push_status(path, fetch))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
if !list.is_empty() {
|
if !list.is_empty() {
|
||||||
@@ -675,14 +677,14 @@ fn compute_statuses(
|
|||||||
}
|
}
|
||||||
|
|
||||||
(statuses, push_statuses)
|
(statuses, push_statuses)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn pr_proposes_checkout(
|
pub fn pr_proposes_checkout(
|
||||||
pr: &Event,
|
pr: &Event,
|
||||||
open: bool,
|
open: bool,
|
||||||
user: PublicKey,
|
user: PublicKey,
|
||||||
checkout: &CheckoutStatus,
|
checkout: &CheckoutStatus,
|
||||||
) -> bool {
|
) -> bool {
|
||||||
if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user {
|
if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
@@ -703,6 +705,7 @@ pub fn pr_proposes_checkout(
|
|||||||
.is_some_and(|tip| tip == checkout.head);
|
.is_some_and(|tip| tip == checkout.head);
|
||||||
|
|
||||||
branch_matches || tip_matches
|
branch_matches || tip_matches
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -763,7 +766,7 @@ mod tests {
|
|||||||
run(&["checkout", "-b", "feature"]);
|
run(&["checkout", "-b", "feature"]);
|
||||||
std::fs::write(path.join("feature.txt"), "x\n").expect("write");
|
std::fs::write(path.join("feature.txt"), "x\n").expect("write");
|
||||||
commit("feature work");
|
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.branch, "feature");
|
||||||
assert_eq!(status.base, "main");
|
assert_eq!(status.base, "main");
|
||||||
assert_eq!(status.ahead, 1);
|
assert_eq!(status.ahead, 1);
|
||||||
@@ -771,12 +774,12 @@ mod tests {
|
|||||||
|
|
||||||
// Dirty worktrees are never suggested.
|
// Dirty worktrees are never suggested.
|
||||||
std::fs::write(path.join("uncommitted.txt"), "y\n").expect("write");
|
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", "--", "."]);
|
run(&["checkout", "--", "."]);
|
||||||
|
|
||||||
// Even on main, nothing to propose.
|
// Even on main, nothing to propose.
|
||||||
run(&["checkout", "main"]);
|
run(&["checkout", "main"]);
|
||||||
assert_eq!(checkout_status(&path, Some("main")), None);
|
assert_eq!(CheckoutsStore::checkout_status(&path, Some("main")), None);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -822,13 +825,13 @@ mod tests {
|
|||||||
};
|
};
|
||||||
|
|
||||||
// A fresh clone has nothing to push.
|
// 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.
|
// One local commit, ready to push, counted against the remote.
|
||||||
std::fs::write(checkout.join("work.txt"), "x\n").expect("write");
|
std::fs::write(checkout.join("work.txt"), "x\n").expect("write");
|
||||||
run(&["add", "-A"]);
|
run(&["add", "-A"]);
|
||||||
run(&["commit", "-m", "local work"]);
|
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.branch, "main");
|
||||||
assert_eq!(status.base, "refs/remotes/origin/main");
|
assert_eq!(status.base, "refs/remotes/origin/main");
|
||||||
assert_eq!(status.ahead, 1);
|
assert_eq!(status.ahead, 1);
|
||||||
@@ -836,12 +839,12 @@ mod tests {
|
|||||||
|
|
||||||
// The local-only pass reads the tracking refs, no fetch needed:
|
// The local-only pass reads the tracking refs, no fetch needed:
|
||||||
// a commit lands locally long before the remote is reconciled.
|
// 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);
|
assert_eq!(local.ahead, 1);
|
||||||
|
|
||||||
// After the push the same commit is on the remote, idle again.
|
// After the push the same commit is on the remote, idle again.
|
||||||
run(&["push", "origin", "main"]);
|
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.
|
// A commit made by someone else on the remote must not count as local work.
|
||||||
// It is behind, not ahead.
|
// It is behind, not ahead.
|
||||||
@@ -860,6 +863,6 @@ mod tests {
|
|||||||
std::fs::write(remote.join("other.txt"), "y\n").expect("write");
|
std::fs::write(remote.join("other.txt"), "y\n").expect("write");
|
||||||
remote_run(&["add", "-A"]);
|
remote_run(&["add", "-A"]);
|
||||||
remote_run(&["commit", "-m", "remote work"]);
|
remote_run(&["commit", "-m", "remote work"]);
|
||||||
assert_eq!(checkout_push_status(&checkout, true), None);
|
assert_eq!(CheckoutsStore::checkout_push_status(&checkout, true), None);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,34 +8,40 @@ use signed_git::{GitCache, Repo};
|
|||||||
|
|
||||||
static GIT_CACHE: OnceLock<GitCache> = OnceLock::new();
|
static GIT_CACHE: OnceLock<GitCache> = OnceLock::new();
|
||||||
|
|
||||||
fn git_cache() -> &'static GitCache {
|
/// The global git repository mirror cache.
|
||||||
GIT_CACHE
|
pub struct Mirrors;
|
||||||
.get()
|
|
||||||
.expect("git cache is initialized by signed_state::init")
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The root directory of the repository mirrors.
|
impl Mirrors {
|
||||||
pub(crate) fn repo_mirror_root() -> PathBuf {
|
/// Install the global mirror cache root, once.
|
||||||
git_cache().root().to_path_buf()
|
pub fn install(root: impl Into<PathBuf>) {
|
||||||
}
|
|
||||||
|
|
||||||
/// The on-disk path of the mirror of `addr`.
|
|
||||||
pub fn repo_mirror_path(addr: &RepoAddr) -> PathBuf {
|
|
||||||
git_cache().repo_path(addr)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Open the mirror of `addr`, if it has been cloned.
|
|
||||||
pub fn open_repo_mirror(addr: &RepoAddr) -> Result<Option<Repo>> {
|
|
||||||
git_cache().open(addr)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Open the mirror of `addr`, cloning it first when it does not exist yet.
|
|
||||||
pub fn ensure_repo_mirror(addr: &RepoAddr, clone_urls: &[Url]) -> Result<Repo> {
|
|
||||||
git_cache().ensure_clone(addr, clone_urls)
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn set_git_cache(root: impl Into<PathBuf>) {
|
|
||||||
if GIT_CACHE.set(GitCache::new(root.into())).is_err() {
|
if GIT_CACHE.set(GitCache::new(root.into())).is_err() {
|
||||||
log::warn!("git cache root is already set, keeping the first one");
|
log::warn!("git cache root is already set, keeping the first one");
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn cache() -> &'static GitCache {
|
||||||
|
GIT_CACHE
|
||||||
|
.get()
|
||||||
|
.expect("git cache is initialized by signed_state::init")
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The root directory of the repository mirrors.
|
||||||
|
pub(crate) fn root() -> PathBuf {
|
||||||
|
Self::cache().root().to_path_buf()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The on-disk path of the mirror of `addr`.
|
||||||
|
pub fn path(addr: &RepoAddr) -> PathBuf {
|
||||||
|
Self::cache().repo_path(addr)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Open the mirror of `addr`, if it has been cloned.
|
||||||
|
pub fn open(addr: &RepoAddr) -> Result<Option<Repo>> {
|
||||||
|
Self::cache().open(addr)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Open the mirror of `addr`, cloning it first when it does not exist yet.
|
||||||
|
pub fn ensure(addr: &RepoAddr, clone_urls: &[Url]) -> Result<Repo> {
|
||||||
|
Self::cache().ensure_clone(addr, clone_urls)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -36,7 +36,7 @@ impl Inbox {
|
|||||||
cx.notify();
|
cx.notify();
|
||||||
|
|
||||||
let backend = Backend::global(cx);
|
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| {
|
cx.spawn(async move |this, cx| {
|
||||||
let loaded = work.await;
|
let loaded = work.await;
|
||||||
@@ -82,7 +82,7 @@ impl Inbox {
|
|||||||
let state = self.state.clone();
|
let state = self.state.clone();
|
||||||
|
|
||||||
let task: Task<Result<(), Error>> = cx.background_spawn(async move {
|
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}");
|
log::warn!("failed to save inbox state: {error}");
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -90,17 +90,18 @@ impl Inbox {
|
|||||||
|
|
||||||
task.detach();
|
task.detach();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn query_inbox(
|
/// The notifications and authored activity of `me`, grouped into inbox items.
|
||||||
|
pub async fn query(
|
||||||
client: &Client,
|
client: &Client,
|
||||||
me: PublicKey,
|
me: PublicKey,
|
||||||
state: &InboxReadState,
|
state: &InboxReadState,
|
||||||
) -> Result<(Vec<InboxItem>, usize), Error> {
|
) -> Result<(Vec<InboxItem>, usize), Error> {
|
||||||
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);
|
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();
|
let mut activity = Vec::new();
|
||||||
for event in client
|
for event in client
|
||||||
@@ -122,17 +123,19 @@ pub async fn query_inbox(
|
|||||||
let unread_count = items.iter().filter(|item| item.is_unread()).count();
|
let unread_count = items.iter().filter(|item| item.is_unread()).count();
|
||||||
|
|
||||||
Ok((items, unread_count))
|
Ok((items, unread_count))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `d` tag identifying the inbox state event of `me`.
|
impl Inbox {
|
||||||
fn inbox_state_d_tag(me: PublicKey) -> String {
|
/// `d` tag identifying the inbox state event of `me`.
|
||||||
|
fn inbox_state_d_tag(me: PublicKey) -> String {
|
||||||
format!("signed-inbox-state:{}", me.to_hex())
|
format!("signed-inbox-state:{}", me.to_hex())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn load_state(client: &Client, me: PublicKey) -> Result<Option<InboxReadState>, Error> {
|
async fn load_state(client: &Client, me: PublicKey) -> Result<Option<InboxReadState>, Error> {
|
||||||
let filter = Filter::new()
|
let filter = Filter::new()
|
||||||
.kind(Kind::ApplicationSpecificData)
|
.kind(Kind::ApplicationSpecificData)
|
||||||
.identifier(inbox_state_d_tag(me));
|
.identifier(Self::inbox_state_d_tag(me));
|
||||||
|
|
||||||
let events = client.database().query(filter).await?;
|
let events = client.database().query(filter).await?;
|
||||||
|
|
||||||
@@ -147,25 +150,29 @@ async fn load_state(client: &Client, me: PublicKey) -> Result<Option<InboxReadSt
|
|||||||
Ok(None)
|
Ok(None)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Sign with a random key and store locally.
|
/// Sign with a random key and store locally.
|
||||||
async fn save_state(client: &Client, me: PublicKey, state: &InboxReadState) -> Result<(), Error> {
|
async fn save_state(
|
||||||
|
client: &Client,
|
||||||
|
me: PublicKey,
|
||||||
|
state: &InboxReadState,
|
||||||
|
) -> Result<(), Error> {
|
||||||
let event = EventBuilder::new(Kind::ApplicationSpecificData, serde_json::to_string(state)?)
|
let event = EventBuilder::new(Kind::ApplicationSpecificData, serde_json::to_string(state)?)
|
||||||
.tags([Tag::identifier(inbox_state_d_tag(me))])
|
.tags([Tag::identifier(Self::inbox_state_d_tag(me))])
|
||||||
.finalize(&Keys::generate())?;
|
.finalize(&Keys::generate())?;
|
||||||
|
|
||||||
client.database().save_event(&event).await?;
|
client.database().save_event(&event).await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Notification events and a lookup of every ancestor they reference.
|
/// Notification events and a lookup of every ancestor they reference.
|
||||||
async fn fetch_notifications(
|
async fn fetch_notifications(
|
||||||
client: &Client,
|
client: &Client,
|
||||||
me: PublicKey,
|
me: PublicKey,
|
||||||
deletions: &Deletions,
|
deletions: &Deletions,
|
||||||
) -> Result<(Vec<Event>, HashMap<EventId, Event>), Error> {
|
) -> Result<(Vec<Event>, HashMap<EventId, Event>), Error> {
|
||||||
let mut notifications: Vec<Event> = Vec::new();
|
let mut notifications: Vec<Event> = Vec::new();
|
||||||
let mut by_id: HashMap<EventId, Event> = HashMap::new();
|
let mut by_id: HashMap<EventId, Event> = HashMap::new();
|
||||||
|
|
||||||
@@ -181,7 +188,11 @@ async fn fetch_notifications(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut pending: Vec<EventId> = notifications.iter().flat_map(event_references).collect();
|
let mut pending: Vec<EventId> = notifications
|
||||||
|
.iter()
|
||||||
|
.flat_map(Self::event_references)
|
||||||
|
.collect();
|
||||||
|
|
||||||
let mut seen: HashSet<EventId> = by_id.keys().copied().collect();
|
let mut seen: HashSet<EventId> = by_id.keys().copied().collect();
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
@@ -202,7 +213,7 @@ async fn fetch_notifications(
|
|||||||
if deletions.is_deleted(&event) {
|
if deletions.is_deleted(&event) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
next.extend(event_references(&event).filter(|id| !seen.contains(id)));
|
next.extend(Self::event_references(&event).filter(|id| !seen.contains(id)));
|
||||||
by_id.entry(event.id).or_insert(event);
|
by_id.entry(event.id).or_insert(event);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -210,10 +221,10 @@ async fn fetch_notifications(
|
|||||||
}
|
}
|
||||||
|
|
||||||
Ok((notifications, by_id))
|
Ok((notifications, by_id))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Event ids referenced by `event` through its `e` and `E` tags.
|
/// Event ids referenced by `event` through its `e` and `E` tags.
|
||||||
fn event_references(event: &Event) -> impl Iterator<Item = EventId> + '_ {
|
fn event_references(event: &Event) -> impl Iterator<Item = EventId> + '_ {
|
||||||
event.tags.iter().filter_map(|tag| {
|
event.tags.iter().filter_map(|tag| {
|
||||||
if tag.kind() != "e" && tag.kind() != "E" {
|
if tag.kind() != "e" && tag.kind() != "E" {
|
||||||
return None;
|
return None;
|
||||||
@@ -221,4 +232,5 @@ fn event_references(event: &Event) -> impl Iterator<Item = EventId> + '_ {
|
|||||||
tag.content()
|
tag.content()
|
||||||
.and_then(|content| EventId::from_hex(content).ok())
|
.and_then(|content| EventId::from_hex(content).ok())
|
||||||
})
|
})
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,24 +1,27 @@
|
|||||||
mod backend;
|
mod backend;
|
||||||
|
mod bootstrap;
|
||||||
mod checkouts;
|
mod checkouts;
|
||||||
mod git_store;
|
mod git_store;
|
||||||
mod inbox;
|
mod inbox;
|
||||||
mod local_repos;
|
mod local_repos;
|
||||||
mod profile;
|
mod profile;
|
||||||
|
mod push;
|
||||||
mod refresh;
|
mod refresh;
|
||||||
mod repo;
|
mod repo;
|
||||||
mod repos;
|
mod repos;
|
||||||
|
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
|
|
||||||
pub use backend::{Backend, BackendEvent, user_grasp_list_servers};
|
pub use backend::{Backend, BackendEvent};
|
||||||
pub use checkouts::{CheckoutStatus, CheckoutsStore, pr_proposes_checkout};
|
pub use bootstrap::user_grasp_list_servers;
|
||||||
use git_store::set_git_cache;
|
pub use checkouts::{CheckoutStatus, CheckoutsStore};
|
||||||
pub use git_store::{ensure_repo_mirror, open_repo_mirror, repo_mirror_path};
|
pub use git_store::Mirrors;
|
||||||
use gpui::{App, AppContext};
|
use gpui::{App, AppContext};
|
||||||
pub use inbox::{Inbox, query_inbox};
|
pub use inbox::Inbox;
|
||||||
pub use local_repos::{LocalReposStore, ResolvedLocalRepo, resolve_local_repos};
|
pub use local_repos::{LocalReposStore, ResolvedLocalRepo};
|
||||||
pub use nostr_sdk::prelude::Timestamp;
|
pub use nostr_sdk::prelude::Timestamp;
|
||||||
pub use profile::{Profile, ProfileStore};
|
pub use profile::{Profile, ProfileStore};
|
||||||
|
pub use push::{GraspServer, PushOutcome};
|
||||||
pub use refresh::{RefreshGate, RefreshRequest};
|
pub use refresh::{RefreshGate, RefreshRequest};
|
||||||
pub use repo::RepoStore;
|
pub use repo::RepoStore;
|
||||||
pub use repos::RepoListStore;
|
pub use repos::RepoListStore;
|
||||||
@@ -42,7 +45,7 @@ pub fn init(
|
|||||||
let (client, signer) = (backend.client, backend.signer);
|
let (client, signer) = (backend.client, backend.signer);
|
||||||
|
|
||||||
// Set Git cache for the repos root
|
// Set Git cache for the repos root
|
||||||
set_git_cache(repos_root);
|
Mirrors::install(repos_root);
|
||||||
|
|
||||||
// Set global stores for the backend
|
// Set global stores for the backend
|
||||||
Backend::set_global(cx.new(|cx| Backend::new(client, signer, cx)), cx);
|
Backend::set_global(cx.new(|cx| Backend::new(client, signer, cx)), cx);
|
||||||
|
|||||||
@@ -137,11 +137,12 @@ impl ResolvedLocalRepo {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Resolve the scanned repositories against the known announcements.
|
/// Resolve the scanned repositories against the known announcements.
|
||||||
pub fn resolve_local_repos(
|
impl LocalReposStore {
|
||||||
|
pub fn resolve(
|
||||||
repos: &[LocalRepo],
|
repos: &[LocalRepo],
|
||||||
known: &[Announcement],
|
known: &[Announcement],
|
||||||
own: &[Announcement],
|
own: &[Announcement],
|
||||||
) -> Vec<ResolvedLocalRepo> {
|
) -> Vec<ResolvedLocalRepo> {
|
||||||
let shown: HashSet<RepoAddr> = own.iter().map(Announcement::addr).collect();
|
let shown: HashSet<RepoAddr> = own.iter().map(Announcement::addr).collect();
|
||||||
|
|
||||||
repos
|
repos
|
||||||
@@ -171,6 +172,7 @@ pub fn resolve_local_repos(
|
|||||||
})
|
})
|
||||||
})
|
})
|
||||||
.collect()
|
.collect()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -222,7 +224,7 @@ mod tests {
|
|||||||
let repo = bound(KEY, "mine");
|
let repo = bound(KEY, "mine");
|
||||||
|
|
||||||
let own = std::slice::from_ref(&own);
|
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]
|
#[test]
|
||||||
@@ -230,7 +232,7 @@ mod tests {
|
|||||||
let known = announcement(OTHER_KEY, "theirs");
|
let known = announcement(OTHER_KEY, "theirs");
|
||||||
let repo = bound(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.len(), 1);
|
||||||
assert_eq!(resolved[0].announcement.as_ref(), Some(&known));
|
assert_eq!(resolved[0].announcement.as_ref(), Some(&known));
|
||||||
|
|||||||
@@ -10,7 +10,8 @@ use gpui::{
|
|||||||
use nostr_sdk::prelude::*;
|
use nostr_sdk::prelude::*;
|
||||||
use utils::shorten_pubkey;
|
use utils::shorten_pubkey;
|
||||||
|
|
||||||
use crate::backend::{Backend, BackendEvent, sync_bootstrap_only};
|
use crate::backend::{Backend, BackendEvent};
|
||||||
|
use crate::bootstrap::sync_bootstrap_only;
|
||||||
|
|
||||||
/// How long to wait for more requests before firing a batched fetch.
|
/// How long to wait for more requests before firing a batched fetch.
|
||||||
const BATCH_TIMEOUT: Duration = Duration::from_millis(500);
|
const BATCH_TIMEOUT: Duration = Duration::from_millis(500);
|
||||||
|
|||||||
@@ -0,0 +1,577 @@
|
|||||||
|
use std::path::Path;
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
use anyhow::{Error, bail};
|
||||||
|
use gpui::BackgroundExecutor;
|
||||||
|
use nostr::event::IntoEventBuilder;
|
||||||
|
use nostr::prelude::Url;
|
||||||
|
use nostr_sdk::prelude::*;
|
||||||
|
use signed_core::RepoState;
|
||||||
|
use signed_git::Repo;
|
||||||
|
use signed_nostr::UniversalSigner;
|
||||||
|
|
||||||
|
pub(crate) const GRASP_PUSH_ATTEMPTS: usize = 3;
|
||||||
|
|
||||||
|
/// Pause before re-staging a state event after a transient denial.
|
||||||
|
const GRASP_RETRY_DELAY: Duration = Duration::from_secs(1);
|
||||||
|
|
||||||
|
/// Base URL of a grasp server, `https://<host>`.
|
||||||
|
///
|
||||||
|
/// `ws://` grasp servers use `http://<host>`, like ngit.
|
||||||
|
pub(crate) fn grasp_base_url(relay: &RelayUrl) -> Option<String> {
|
||||||
|
// `domain()` drops the port.
|
||||||
|
let parsed = Url::parse(relay.as_str()).ok()?;
|
||||||
|
let host = parsed.host_str()?;
|
||||||
|
let port = parsed.port().map(|p| format!(":{p}")).unwrap_or_default();
|
||||||
|
// `ws://` grasp servers, e.g. local dev relays, speak plain HTTP.
|
||||||
|
let scheme = if relay.scheme().is_secure() {
|
||||||
|
"https"
|
||||||
|
} else {
|
||||||
|
"http"
|
||||||
|
};
|
||||||
|
Some(format!("{scheme}://{host}{port}"))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn grasp_clone_url(relay: &RelayUrl, owner: &str, repo_id: &str) -> Option<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")
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Assemble the `clone` URLs of a pull request.
|
||||||
|
///
|
||||||
|
/// The author's GRASP-06 `/prs/` URLs come first.
|
||||||
|
pub(crate) fn pr_clone_urls(prs_urls: Vec<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,
|
||||||
|
/// `None` when the server accepted the data, the reason otherwise.
|
||||||
|
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
|
||||||
|
}
|
||||||
|
|
||||||
|
/// `None` when the server accepted the data, the reason otherwise.
|
||||||
|
pub fn reason(&self) -> Option<&str> {
|
||||||
|
self.reason.as_deref()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Default)]
|
||||||
|
pub struct PushOutcome {
|
||||||
|
/// Per-server results, in the order the servers were listed.
|
||||||
|
pub servers: Vec<GraspServer>,
|
||||||
|
/// The newest state event a grasp relay accepted for this push, if any.
|
||||||
|
///
|
||||||
|
/// Broadcast to the other relays once a git server holds the data.
|
||||||
|
pub state_event: Option<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("; ")
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A warning for a push only some grasp servers accepted.
|
||||||
|
///
|
||||||
|
/// `None` when every server accepted the push or nothing was pushed.
|
||||||
|
pub fn partial_warning(&self) -> Option<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")
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Keep `event` as the push's fan-out state event when it is newer than the
|
||||||
|
/// current one. All staged events carry the same refs; the newest timestamp
|
||||||
|
/// wins on the relays.
|
||||||
|
fn keep_newest(state_event: &mut Option<Event>, event: Event) {
|
||||||
|
if state_event
|
||||||
|
.as_ref()
|
||||||
|
.is_none_or(|current| event.created_at > current.created_at)
|
||||||
|
{
|
||||||
|
*state_event = Some(event);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The grasp push pipeline: stage a signed state event on each server's
|
||||||
|
/// relay, then push the git data, retrying transient denials.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub(crate) struct GraspPush {
|
||||||
|
client: Client,
|
||||||
|
signer: UniversalSigner,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl GraspPush {
|
||||||
|
pub(crate) fn new(client: Client, signer: UniversalSigner) -> Self {
|
||||||
|
Self { client, signer }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The event was accepted by at least one relay, or a descriptive error otherwise.
|
||||||
|
pub(crate) fn require_relay_accepted(
|
||||||
|
output: SendEventOutput,
|
||||||
|
event: Event,
|
||||||
|
) -> Result<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)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Sign and broadcast `builder`, logging rather than surfacing failures.
|
||||||
|
///
|
||||||
|
/// Used for best-effort identity bootstrap events, where a relay hiccup
|
||||||
|
/// should not block sign-up.
|
||||||
|
pub(crate) async fn publish_best_effort(&self, builder: EventBuilder) {
|
||||||
|
let result: Result<(), Error> = async {
|
||||||
|
let event = builder.finalize_async(&self.signer).await?;
|
||||||
|
let output = self.client.send_event(&event).broadcast().await?;
|
||||||
|
Self::require_relay_accepted(output, event)?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
.await;
|
||||||
|
|
||||||
|
if let Err(e) = result {
|
||||||
|
log::warn!("failed to publish identity bootstrap event: {e}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Sign `builder`, broadcast the event and require a relay to accept it.
|
||||||
|
/// Returns the signed event.
|
||||||
|
pub(crate) async fn publish_one(&self, builder: EventBuilder) -> Result<Event, Error> {
|
||||||
|
let event = builder.finalize_async(&self.signer).await?;
|
||||||
|
self.send_accepted(event).await
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Broadcast an already signed event and require a relay to accept it.
|
||||||
|
/// Returns the event.
|
||||||
|
pub(crate) async fn send_accepted(&self, event: Event) -> Result<Event, Error> {
|
||||||
|
let output = self.client.send_event(&event).broadcast().await?;
|
||||||
|
Self::require_relay_accepted(output, event)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Sign and send a single NIP-09 deletion request for `event`.
|
||||||
|
pub(crate) async fn retract_event(&self, event: &Event) -> Result<(), Error> {
|
||||||
|
let builder = EventDeletionRequest::new()
|
||||||
|
.id(event.id)
|
||||||
|
.into_event_builder();
|
||||||
|
|
||||||
|
let deletion = builder.finalize_async(&self.signer).await?;
|
||||||
|
self.client.send_event(&deletion).broadcast().await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Sign a fresh kind `30618` state event for the push.
|
||||||
|
///
|
||||||
|
/// `last_created_at` is the timestamp of the previous event signed for this push.
|
||||||
|
/// Retries within the same second get the next second: a grasp relay
|
||||||
|
/// treats a same-id resend as a duplicate and does not re-run its ingest,
|
||||||
|
/// so an identical resend cannot re-park a state event lost from its purgatory.
|
||||||
|
async fn sign_state_event(
|
||||||
|
&self,
|
||||||
|
repo_id: &str,
|
||||||
|
refs: &[(String, String)],
|
||||||
|
head: Option<&str>,
|
||||||
|
last_created_at: u64,
|
||||||
|
) -> Result<(Event, u64), String> {
|
||||||
|
let now = Timestamp::now().as_secs();
|
||||||
|
let created_at = if now > last_created_at {
|
||||||
|
now
|
||||||
|
} else {
|
||||||
|
last_created_at + 1
|
||||||
|
};
|
||||||
|
|
||||||
|
let event = RepoState::build(repo_id, refs, head)
|
||||||
|
.custom_created_at(Timestamp::from_secs(created_at))
|
||||||
|
.finalize_async(&self.signer)
|
||||||
|
.await
|
||||||
|
.map_err(|e| format!("could not sign the state event: {e}"))?;
|
||||||
|
|
||||||
|
Ok((event, created_at))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Ensure the relay is known and connected, then publish `event` to it.
|
||||||
|
///
|
||||||
|
/// `Ok` only when the relay confirmed the event.
|
||||||
|
/// On a grasp relay the accept parks the event in purgatory,
|
||||||
|
/// which authorizes the paired git push.
|
||||||
|
async fn stage_event_on_relay(&self, relay: &RelayUrl, event: &Event) -> Result<(), String> {
|
||||||
|
self.client
|
||||||
|
.add_relay(relay)
|
||||||
|
.and_connect()
|
||||||
|
.await
|
||||||
|
.map_err(|e| format!("could not add relay {relay}: {e}"))?;
|
||||||
|
|
||||||
|
let output = self
|
||||||
|
.client
|
||||||
|
.send_event(event)
|
||||||
|
.to([relay.clone()])
|
||||||
|
.await
|
||||||
|
.map_err(|e| format!("could not send the state event to {relay}: {e}"))?;
|
||||||
|
|
||||||
|
if output.success.contains_key(relay) {
|
||||||
|
Ok(())
|
||||||
|
} else {
|
||||||
|
let reason = output
|
||||||
|
.failed
|
||||||
|
.get(relay)
|
||||||
|
.cloned()
|
||||||
|
.unwrap_or_else(|| "relay did not confirm the event".to_owned());
|
||||||
|
Err(reason)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[allow(clippy::too_many_arguments)]
|
||||||
|
pub(crate) async fn push_staged_to_grasps(
|
||||||
|
&self,
|
||||||
|
repo_id: &str,
|
||||||
|
refs: &[(String, String)],
|
||||||
|
head: Option<&str>,
|
||||||
|
path: &Path,
|
||||||
|
owner: &str,
|
||||||
|
servers: &[RelayUrl],
|
||||||
|
executor: &BackgroundExecutor,
|
||||||
|
push: impl Fn(&Path, &str, &str, &str) -> Result<(), Error>,
|
||||||
|
) -> PushOutcome {
|
||||||
|
let mut outcome = PushOutcome::default();
|
||||||
|
|
||||||
|
for relay in servers {
|
||||||
|
let Some(base) = grasp_base_url(relay) else {
|
||||||
|
outcome
|
||||||
|
.servers
|
||||||
|
.push(GraspServer::failed(relay.clone(), "no domain"));
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
let git_url = format!("{base}/{owner}/{repo_id}.git");
|
||||||
|
|
||||||
|
let mut reason = None;
|
||||||
|
let mut last_created_at = 0;
|
||||||
|
// The last state event staged on this server, for the convergence
|
||||||
|
// probe below when every push attempt lost the stale-ref race.
|
||||||
|
let mut staged_event = None;
|
||||||
|
|
||||||
|
'server: for attempt in 1..=GRASP_PUSH_ATTEMPTS {
|
||||||
|
if attempt > 1 {
|
||||||
|
// Give the server's ingest a moment before re-staging.
|
||||||
|
executor.timer(GRASP_RETRY_DELAY).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
let (event, created_at) = match self
|
||||||
|
.sign_state_event(repo_id, refs, head, last_created_at)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(signed) => signed,
|
||||||
|
Err(e) => {
|
||||||
|
reason = Some(e);
|
||||||
|
break 'server;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
last_created_at = created_at;
|
||||||
|
|
||||||
|
// Stage the state event on this server's own relay.
|
||||||
|
// A failed stage means the grasp never parked the state,
|
||||||
|
// so the git push would be denied anyway: skip it (the eligibility gate).
|
||||||
|
if let Err(e) = self.stage_event_on_relay(relay, &event).await {
|
||||||
|
// One retry absorbs a relay connect blip, on the first
|
||||||
|
// attempt only.
|
||||||
|
if attempt == 1 && self.stage_event_on_relay(relay, &event).await.is_ok() {
|
||||||
|
// staged on the retry
|
||||||
|
} else {
|
||||||
|
reason = Some(e);
|
||||||
|
break 'server;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
staged_event = Some(event.clone());
|
||||||
|
|
||||||
|
match push(path, &base, owner, repo_id) {
|
||||||
|
Ok(()) => {
|
||||||
|
keep_newest(&mut outcome.state_event, event);
|
||||||
|
break 'server;
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
let text = e.to_string();
|
||||||
|
if attempt < GRASP_PUSH_ATTEMPTS && is_transient_grasp_denial(&text) {
|
||||||
|
reason = Some(text);
|
||||||
|
continue 'server;
|
||||||
|
}
|
||||||
|
reason = Some(text);
|
||||||
|
break 'server;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The grasp's own background sync aligns refs to staged state
|
||||||
|
// events as soon as the objects land, which can beat every push
|
||||||
|
// attempt's compare-and-swap (`cannot lock ref ... but expected`).
|
||||||
|
// When the last denial was that race the sync has usually finished
|
||||||
|
// by now: verify the advertised refs and accept the server when the
|
||||||
|
// pushed data is already there.
|
||||||
|
if let Some(last_reason) = &reason
|
||||||
|
&& is_stale_advertisement_race(last_reason)
|
||||||
|
&& Repo::open(path)
|
||||||
|
.and_then(|repo| repo.remote_has_refs(&git_url, refs))
|
||||||
|
.unwrap_or(false)
|
||||||
|
{
|
||||||
|
if let Some(event) = staged_event {
|
||||||
|
keep_newest(&mut outcome.state_event, event);
|
||||||
|
}
|
||||||
|
reason = None;
|
||||||
|
}
|
||||||
|
|
||||||
|
match reason {
|
||||||
|
Some(reason) => {
|
||||||
|
log::warn!("grasp push failed: {relay}: {reason}");
|
||||||
|
outcome
|
||||||
|
.servers
|
||||||
|
.push(GraspServer::failed(relay.clone(), reason));
|
||||||
|
}
|
||||||
|
None => outcome.servers.push(GraspServer::ok(relay.clone())),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
outcome
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn grasp_base_url_maps_schemes_like_ngit() {
|
||||||
|
let wss = RelayUrl::parse("wss://relay.ngit.dev").expect("url");
|
||||||
|
assert_eq!(
|
||||||
|
grasp_base_url(&wss).as_deref(),
|
||||||
|
Some("https://relay.ngit.dev")
|
||||||
|
);
|
||||||
|
|
||||||
|
let ws = RelayUrl::parse("ws://localhost:8080").expect("url");
|
||||||
|
assert_eq!(
|
||||||
|
grasp_base_url(&ws).as_deref(),
|
||||||
|
Some("http://localhost:8080")
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn grasp_clone_url_matches_ngit_format() {
|
||||||
|
let relay = RelayUrl::parse("wss://gitnostr.com").expect("url");
|
||||||
|
let url = grasp_clone_url(&relay, "npub1test", "my-repo").expect("url");
|
||||||
|
assert_eq!(
|
||||||
|
url.to_string(),
|
||||||
|
"https://gitnostr.com/npub1test/my-repo.git"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn grasp06_prs_url_matches_ngit_format() {
|
||||||
|
assert_eq!(
|
||||||
|
grasp06_prs_url("https://relay.ngit.dev", "npub1author", "my-repo"),
|
||||||
|
"https://relay.ngit.dev/prs/npub1author/my-repo.git"
|
||||||
|
);
|
||||||
|
// `ws://` grasp servers, local dev, keep their plain-HTTP base.
|
||||||
|
assert_eq!(
|
||||||
|
grasp06_prs_url("http://localhost:8080", "npub1author", "my-repo"),
|
||||||
|
"http://localhost:8080/prs/npub1author/my-repo.git"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn pr_clone_urls_orders_author_first_and_deduplicates() {
|
||||||
|
let prs = vec![
|
||||||
|
Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"),
|
||||||
|
Url::parse("https://a.example/prs/npub1me/repo.git").expect("url"),
|
||||||
|
];
|
||||||
|
let base = vec![
|
||||||
|
Url::parse("https://a.example/npub1owner/repo.git").expect("url"),
|
||||||
|
Url::parse("https://b.example/npub1owner/repo.git").expect("url"),
|
||||||
|
Url::parse("https://b.example/npub1owner/repo.git").expect("url"),
|
||||||
|
];
|
||||||
|
|
||||||
|
let urls = pr_clone_urls(prs, base);
|
||||||
|
assert_eq!(
|
||||||
|
urls.iter().map(ToString::to_string).collect::<Vec<_>>(),
|
||||||
|
vec![
|
||||||
|
"https://a.example/prs/npub1me/repo.git",
|
||||||
|
"https://a.example/npub1owner/repo.git",
|
||||||
|
"https://b.example/npub1owner/repo.git",
|
||||||
|
]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn transient_grasp_denials_are_classified() {
|
||||||
|
// The exact server rejection that started this work: the state event
|
||||||
|
// had not reached the grasp's purgatory before the git push.
|
||||||
|
let reported = "remote: ERR authorisation failed: No state events in purgatory\n\
|
||||||
|
fatal: the remote end hung up unexpectedly\n\
|
||||||
|
error: failed to push some refs to 'https://relay.ngit.dev/...git'";
|
||||||
|
assert!(is_transient_grasp_denial(reported));
|
||||||
|
|
||||||
|
// The other purgatory states a fresh event resolves.
|
||||||
|
assert!(is_transient_grasp_denial(
|
||||||
|
"remote: ERR authorisation failed: No matching state event found in purgatory"
|
||||||
|
));
|
||||||
|
assert!(is_transient_grasp_denial(
|
||||||
|
"remote: ERR authorisation failed: 1 state event in purgatory from authorized \
|
||||||
|
publisher but doesn't match push"
|
||||||
|
));
|
||||||
|
assert!(is_transient_grasp_denial(
|
||||||
|
"remote: ERR authorisation failed: 2 state events in purgatory but none from \
|
||||||
|
authorized publishers"
|
||||||
|
));
|
||||||
|
assert!(is_transient_grasp_denial(
|
||||||
|
"remote: ERR authorisation failed: No repository announcement found"
|
||||||
|
));
|
||||||
|
|
||||||
|
// Rejections a fresh state event cannot fix are not retried.
|
||||||
|
assert!(!is_transient_grasp_denial(
|
||||||
|
"remote: ERR authorisation failed: not a maintainer of this repository"
|
||||||
|
));
|
||||||
|
assert!(!is_transient_grasp_denial(
|
||||||
|
"fatal: unable to access 'https://relay.ngit.dev/...': The requested URL returned \
|
||||||
|
error: 403"
|
||||||
|
));
|
||||||
|
assert!(!is_transient_grasp_denial(
|
||||||
|
"fatal: unable to access 'https://relay.ngit.dev/...': Could not resolve host"
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn stale_ref_races_are_retried() {
|
||||||
|
// The grasp's background sync aligned the ref to a parked state event
|
||||||
|
// between this push's advertisement and its ref transaction. The ref
|
||||||
|
// is usually already where the push wants it, so a retry converges.
|
||||||
|
let reported = "remote: error: cannot lock ref 'refs/heads/main': is at \
|
||||||
|
cac2ac91b6f5fb8dfcb6962785babc6e65350cb3 but expected \
|
||||||
|
bc5e892aa84dc6240a5fbcd59367a4857d26f49b\n\
|
||||||
|
To https://relay.ngit.dev/npub1owner/signed-test.git\n\
|
||||||
|
! [remote rejected] main -> main (incorrect old value provided)\n\
|
||||||
|
error: failed to push some refs to 'https://relay.ngit.dev/npub1owner/signed-test.git'";
|
||||||
|
assert!(is_transient_grasp_denial(reported));
|
||||||
|
assert!(is_stale_advertisement_race(reported));
|
||||||
|
|
||||||
|
// Markers match independently of the surrounding git output.
|
||||||
|
assert!(is_stale_advertisement_race(
|
||||||
|
"cannot lock ref 'refs/heads/main'"
|
||||||
|
));
|
||||||
|
assert!(is_stale_advertisement_race(
|
||||||
|
"! [remote rejected] main -> main (incorrect old value provided)"
|
||||||
|
));
|
||||||
|
|
||||||
|
// A purgatory denial is not a stale-advertisement race.
|
||||||
|
assert!(!is_stale_advertisement_race("No state events in purgatory"));
|
||||||
|
|
||||||
|
// A real divergence is a different error and stays permanent.
|
||||||
|
assert!(!is_transient_grasp_denial(
|
||||||
|
" ! [rejected] main -> main (non-fast-forward)"
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -15,11 +15,10 @@ use signed_core::{
|
|||||||
use signed_git::{Nip34Binding, PatchParser, Repo};
|
use signed_git::{Nip34Binding, PatchParser, Repo};
|
||||||
use signed_nostr::UniversalSigner;
|
use signed_nostr::UniversalSigner;
|
||||||
|
|
||||||
use crate::backend::{
|
use crate::backend::{Backend, BackendEvent};
|
||||||
Backend, BackendEvent, grasp_base_url, grasp06_prs_url, pr_clone_urls, require_relay_accepted,
|
use crate::bootstrap::user_grasp_list_servers;
|
||||||
user_grasp_list_servers,
|
|
||||||
};
|
|
||||||
use crate::checkouts::CheckoutsStore;
|
use crate::checkouts::CheckoutsStore;
|
||||||
|
use crate::push::{GraspPush, PushOutcome, grasp_base_url, grasp06_prs_url, pr_clone_urls};
|
||||||
use crate::repos::RepoListStore;
|
use crate::repos::RepoListStore;
|
||||||
|
|
||||||
/// Maximum size of one patch event.
|
/// Maximum size of one patch event.
|
||||||
@@ -700,19 +699,12 @@ impl RepoStore {
|
|||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
|
|
||||||
let series: Vec<String> = PatchParser::split_patch_series(&patch)
|
let series = PatchSeries::parse(&patch);
|
||||||
.into_iter()
|
|
||||||
.map(str::to_owned)
|
|
||||||
.collect();
|
|
||||||
|
|
||||||
if let Some(oversized) = series
|
if let Some(oversized) = series.oversized_length() {
|
||||||
.iter()
|
|
||||||
.find(|part| part.len() > MAX_PATCH_EVENT_BYTES)
|
|
||||||
{
|
|
||||||
self.last_error = Some(format!(
|
self.last_error = Some(format!(
|
||||||
"patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)",
|
"patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)",
|
||||||
oversized.len(),
|
oversized, MAX_PATCH_EVENT_BYTES
|
||||||
MAX_PATCH_EVENT_BYTES
|
|
||||||
));
|
));
|
||||||
cx.notify();
|
cx.notify();
|
||||||
return;
|
return;
|
||||||
@@ -720,11 +712,7 @@ impl RepoStore {
|
|||||||
|
|
||||||
// The tip of the series is its last commit.
|
// The tip of the series is its last commit.
|
||||||
// `git format-patch` orders patches oldest first.
|
// `git format-patch` orders patches oldest first.
|
||||||
let Some(current_commit) = series
|
let Some(current_commit) = series.tip_commit() else {
|
||||||
.last()
|
|
||||||
.and_then(|part| patch_current_commit(part))
|
|
||||||
.and_then(|hex| hex.parse::<Sha1Hash>().ok())
|
|
||||||
else {
|
|
||||||
self.last_error = Some(
|
self.last_error = Some(
|
||||||
"Patch must be `git format-patch` output with a `From <commit-id>` header".into(),
|
"Patch must be `git format-patch` output with a `From <commit-id>` header".into(),
|
||||||
);
|
);
|
||||||
@@ -930,11 +918,10 @@ impl RepoStore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
let publish_result: Result<Event, Error> = async {
|
let publish_result = {
|
||||||
let output = client.send_event(&event).broadcast().await?;
|
let pusher = GraspPush::new(client.clone(), signer.clone());
|
||||||
require_relay_accepted(output, event)
|
pusher.send_accepted(event).await
|
||||||
}
|
};
|
||||||
.await;
|
|
||||||
|
|
||||||
let pr_event = match publish_result {
|
let pr_event = match publish_result {
|
||||||
Ok(event) => event,
|
Ok(event) => event,
|
||||||
@@ -1037,29 +1024,18 @@ impl RepoStore {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
let series: Vec<String> = PatchParser::split_patch_series(&patch)
|
let series = PatchSeries::parse(&patch);
|
||||||
.into_iter()
|
if let Some(oversized) = series.oversized_length() {
|
||||||
.map(str::to_owned)
|
|
||||||
.collect();
|
|
||||||
if let Some(oversized) = series
|
|
||||||
.iter()
|
|
||||||
.find(|part| part.len() > MAX_PATCH_EVENT_BYTES)
|
|
||||||
{
|
|
||||||
self.last_error = Some(format!(
|
self.last_error = Some(format!(
|
||||||
"patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)",
|
"patch too large ({} bytes; NIP-34 suggests keeping each patch under {} bytes)",
|
||||||
oversized.len(),
|
oversized, MAX_PATCH_EVENT_BYTES
|
||||||
MAX_PATCH_EVENT_BYTES
|
|
||||||
));
|
));
|
||||||
cx.notify();
|
cx.notify();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// The new tip of the PR is the last commit of the series.
|
// The new tip of the PR is the last commit of the series.
|
||||||
let Some(current_commit) = series
|
let Some(current_commit) = series.tip_commit() else {
|
||||||
.last()
|
|
||||||
.and_then(|part| patch_current_commit(part))
|
|
||||||
.and_then(|hex| hex.parse::<Sha1Hash>().ok())
|
|
||||||
else {
|
|
||||||
self.last_error = Some(
|
self.last_error = Some(
|
||||||
"Patch must be `git format-patch` output with a `From <commit-id>` header".into(),
|
"Patch must be `git format-patch` output with a `From <commit-id>` header".into(),
|
||||||
);
|
);
|
||||||
@@ -1129,12 +1105,10 @@ impl RepoStore {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let publish_result: Result<Event, Error> = async {
|
let publish_result = {
|
||||||
let event = builder.finalize_async(&signer).await?;
|
let pusher = GraspPush::new(client.clone(), signer.clone());
|
||||||
let output = client.send_event(&event).broadcast().await?;
|
pusher.publish_one(builder).await
|
||||||
require_relay_accepted(output, event)
|
};
|
||||||
}
|
|
||||||
.await;
|
|
||||||
|
|
||||||
if let Err(e) = publish_result {
|
if let Err(e) = publish_result {
|
||||||
return this.update(cx, |this, cx| {
|
return this.update(cx, |this, cx| {
|
||||||
@@ -1216,38 +1190,10 @@ impl RepoStore {
|
|||||||
return self.action_error("Repository announcement is not loaded yet", cx);
|
return self.action_error("Repository announcement is not loaded yet", cx);
|
||||||
};
|
};
|
||||||
|
|
||||||
self.pushing = true;
|
|
||||||
self.last_error = None;
|
|
||||||
self.last_push_warning = None;
|
|
||||||
cx.notify();
|
|
||||||
|
|
||||||
let backend = Backend::global(cx);
|
let backend = Backend::global(cx);
|
||||||
let push = backend.update(cx, |backend, cx| backend.push_repository(announcement, cx));
|
let push = backend.update(cx, |backend, cx| backend.push_repository(announcement, cx));
|
||||||
|
|
||||||
cx.spawn(async move |this, cx| {
|
self.run_push(push, None, cx)
|
||||||
let result = push.await;
|
|
||||||
|
|
||||||
this.update(cx, |this, cx| {
|
|
||||||
this.pushing = false;
|
|
||||||
|
|
||||||
match &result {
|
|
||||||
Ok(outcome) => {
|
|
||||||
this.last_error = None;
|
|
||||||
// A push only some grasp servers accepted is a warning:
|
|
||||||
// the repo is out of sync on the rest until it is republished.
|
|
||||||
this.last_push_warning = outcome.partial_warning();
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
this.last_error = Some(format!("Push failed: {e}"));
|
|
||||||
this.last_push_warning = None;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
cx.notify();
|
|
||||||
})?;
|
|
||||||
|
|
||||||
result.map(|_| ())
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn push_checkout(
|
pub fn push_checkout(
|
||||||
@@ -1273,17 +1219,30 @@ impl RepoStore {
|
|||||||
// The checkout may be on a side branch.
|
// The checkout may be on a side branch.
|
||||||
let head = self.head.clone();
|
let head = self.head.clone();
|
||||||
|
|
||||||
self.pushing = true;
|
|
||||||
self.last_error = None;
|
|
||||||
self.last_push_warning = None;
|
|
||||||
cx.notify();
|
|
||||||
|
|
||||||
let checkouts = CheckoutsStore::global(cx);
|
|
||||||
let backend = Backend::global(cx);
|
let backend = Backend::global(cx);
|
||||||
let push = backend.update(cx, |backend, cx| {
|
let push = backend.update(cx, |backend, cx| {
|
||||||
backend.push_checkout(announcement, path.clone(), head, cx)
|
backend.push_checkout(announcement, path.clone(), head, cx)
|
||||||
});
|
});
|
||||||
|
|
||||||
|
self.run_push(push, Some((addr, path)), cx)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Run a backend push task, tracking progress in [`Self::pushing`] and
|
||||||
|
/// the outcome in [`Self::last_error`] and [`Self::last_push_warning`].
|
||||||
|
///
|
||||||
|
/// `pushed_checkout` names the checkout whose ready-to-push statuses
|
||||||
|
/// should be recomputed after the remote moved.
|
||||||
|
fn run_push(
|
||||||
|
&mut self,
|
||||||
|
push: Task<Result<PushOutcome, Error>>,
|
||||||
|
pushed_checkout: Option<(RepoAddr, PathBuf)>,
|
||||||
|
cx: &mut Context<Self>,
|
||||||
|
) -> Task<Result<(), Error>> {
|
||||||
|
self.pushing = true;
|
||||||
|
self.last_error = None;
|
||||||
|
self.last_push_warning = None;
|
||||||
|
cx.notify();
|
||||||
|
|
||||||
cx.spawn(async move |this, cx| {
|
cx.spawn(async move |this, cx| {
|
||||||
let result = push.await;
|
let result = push.await;
|
||||||
|
|
||||||
@@ -1296,11 +1255,13 @@ impl RepoStore {
|
|||||||
// A push only some grasp servers accepted is a warning:
|
// A push only some grasp servers accepted is a warning:
|
||||||
// the repo is out of sync on the rest until it is republished.
|
// the repo is out of sync on the rest until it is republished.
|
||||||
this.last_push_warning = outcome.partial_warning();
|
this.last_push_warning = outcome.partial_warning();
|
||||||
|
if let Some((addr, path)) = &pushed_checkout {
|
||||||
// The remote moved, so recompute the ready-to-push statuses.
|
// The remote moved, so recompute the ready-to-push statuses.
|
||||||
checkouts.update(cx, |store, cx| {
|
CheckoutsStore::global(cx).update(cx, |store, cx| {
|
||||||
store.checkout_pushed(&addr, &path, cx);
|
store.checkout_pushed(addr, path, cx);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
this.last_error = Some(format!("Push failed: {e}"));
|
this.last_error = Some(format!("Push failed: {e}"));
|
||||||
this.last_push_warning = None;
|
this.last_push_warning = None;
|
||||||
@@ -1431,12 +1392,8 @@ impl RepoStore {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let task: Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
|
let task: Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
|
||||||
let publish_result: Result<Event, Error> = async {
|
let pusher = GraspPush::new(client, signer);
|
||||||
let event = builder.finalize_async(&signer).await?;
|
let publish_result = pusher.publish_one(builder).await;
|
||||||
let output = client.send_event(&event).broadcast().await?;
|
|
||||||
require_relay_accepted(output, event)
|
|
||||||
}
|
|
||||||
.await;
|
|
||||||
|
|
||||||
if let Err(e) = publish_result {
|
if let Err(e) = publish_result {
|
||||||
this.update(cx, |this, cx| {
|
this.update(cx, |this, cx| {
|
||||||
@@ -1497,6 +1454,46 @@ fn patch_current_commit(patch: &str) -> Option<&str> {
|
|||||||
hex.split_whitespace().next().filter(|hex| hex.len() == 40)
|
hex.split_whitespace().next().filter(|hex| hex.len() == 40)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A `git format-patch` series with the facts derived from its parts.
|
||||||
|
struct PatchSeries {
|
||||||
|
parts: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PatchSeries {
|
||||||
|
fn parse(patch: &str) -> Self {
|
||||||
|
Self {
|
||||||
|
parts: PatchParser::split_patch_series(patch)
|
||||||
|
.into_iter()
|
||||||
|
.map(str::to_owned)
|
||||||
|
.collect(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The byte length of the first part over the NIP-34 size suggestion.
|
||||||
|
fn oversized_length(&self) -> Option<usize> {
|
||||||
|
self.parts
|
||||||
|
.iter()
|
||||||
|
.map(String::len)
|
||||||
|
.find(|length| *length > MAX_PATCH_EVENT_BYTES)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The tip of the series is its last commit;
|
||||||
|
/// `git format-patch` orders patches oldest first.
|
||||||
|
fn tip_commit(&self) -> Option<Sha1Hash> {
|
||||||
|
self.parts
|
||||||
|
.last()
|
||||||
|
.and_then(|part| patch_current_commit(part))
|
||||||
|
.and_then(|hex| hex.parse::<Sha1Hash>().ok())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn commit_of(&self, index: usize) -> Option<&str> {
|
||||||
|
self.parts
|
||||||
|
.get(index)
|
||||||
|
.and_then(|part| patch_current_commit(part))
|
||||||
|
.filter(|hex| hex.len() == 40)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Publish a `git format-patch` series as chained kind-1617 events.
|
/// Publish a `git format-patch` series as chained kind-1617 events.
|
||||||
///
|
///
|
||||||
/// Returns the root event, the one a PR references.
|
/// Returns the root event, the one a PR references.
|
||||||
@@ -1507,15 +1504,16 @@ async fn publish_patch_series(
|
|||||||
addr: &RepoAddr,
|
addr: &RepoAddr,
|
||||||
owner: PublicKey,
|
owner: PublicKey,
|
||||||
euc: Option<&str>,
|
euc: Option<&str>,
|
||||||
series: &[String],
|
series: &PatchSeries,
|
||||||
first_marker: &str,
|
first_marker: &str,
|
||||||
reply_to: Option<EventId>,
|
reply_to: Option<EventId>,
|
||||||
) -> Result<Event, Error> {
|
) -> Result<Event, Error> {
|
||||||
|
let pusher = GraspPush::new(client.clone(), signer.clone());
|
||||||
let mut root: Option<Event> = None;
|
let mut root: Option<Event> = None;
|
||||||
let mut previous = reply_to;
|
let mut previous = reply_to;
|
||||||
|
|
||||||
for (ix, part) in series.iter().enumerate() {
|
for (ix, part) in series.parts.iter().enumerate() {
|
||||||
let Some(commit) = patch_current_commit(part).filter(|hex| hex.len() == 40) else {
|
let Some(commit) = series.commit_of(ix) else {
|
||||||
return Err(anyhow::anyhow!(
|
return Err(anyhow::anyhow!(
|
||||||
"patch {} of the series has no `From <commit-id>` header",
|
"patch {} of the series has no `From <commit-id>` header",
|
||||||
ix + 1
|
ix + 1
|
||||||
@@ -1557,9 +1555,7 @@ async fn publish_patch_series(
|
|||||||
}
|
}
|
||||||
|
|
||||||
let builder = EventBuilder::new(Kind::GitPatch, part.clone()).tags(tags);
|
let builder = EventBuilder::new(Kind::GitPatch, part.clone()).tags(tags);
|
||||||
let event = builder.finalize_async(signer).await?;
|
let event = pusher.publish_one(builder).await?;
|
||||||
let output = client.send_event(&event).broadcast().await?;
|
|
||||||
let event = require_relay_accepted(output, event)?;
|
|
||||||
|
|
||||||
if root.is_none() {
|
if root.is_none() {
|
||||||
root = Some(event.clone());
|
root = Some(event.clone());
|
||||||
|
|||||||
@@ -13,7 +13,7 @@ use gpui_component::{ActiveTheme, Icon, IconName, IconNamed, Sizable, StyledExt,
|
|||||||
use nostr::prelude::{Event, EventId, Kind, PublicKey, Timestamp};
|
use nostr::prelude::{Event, EventId, Kind, PublicKey, Timestamp};
|
||||||
use signed_core::{InboxItem, InboxReadState, RepoAddr};
|
use signed_core::{InboxItem, InboxReadState, RepoAddr};
|
||||||
use signed_state::{
|
use signed_state::{
|
||||||
Backend, BackendEvent, ProfileStore, RefreshGate, RefreshRequest, RepoListStore, query_inbox,
|
Backend, BackendEvent, Inbox, ProfileStore, RefreshGate, RefreshRequest, RepoListStore,
|
||||||
};
|
};
|
||||||
use signed_ui::{Avatar, CountBadge};
|
use signed_ui::{Avatar, CountBadge};
|
||||||
use utils::relative_time;
|
use utils::relative_time;
|
||||||
@@ -199,7 +199,7 @@ impl InboxView {
|
|||||||
|
|
||||||
let state = self.state.clone();
|
let state = self.state.clone();
|
||||||
|
|
||||||
let work = cx.background_spawn(async move { query_inbox(&client, me, &state).await });
|
let work = cx.background_spawn(async move { Inbox::query(&client, me, &state).await });
|
||||||
|
|
||||||
self.tasks.push(cx.spawn(async move |this, cx| {
|
self.tasks.push(cx.spawn(async move |this, cx| {
|
||||||
let (threads, unread_count) = match work.await {
|
let (threads, unread_count) = match work.await {
|
||||||
|
|||||||
@@ -23,7 +23,7 @@ use gpui_component::{
|
|||||||
use nostr::prelude::{Event, EventId, Kind, Url};
|
use nostr::prelude::{Event, EventId, Kind, Url};
|
||||||
use signed_core::{GitEvent, PullRequest, RepoAddr};
|
use signed_core::{GitEvent, PullRequest, RepoAddr};
|
||||||
use signed_git::{FileCommit, PatchParser};
|
use signed_git::{FileCommit, PatchParser};
|
||||||
use signed_state::{Backend, ProfileStore, RepoStore, ensure_repo_mirror};
|
use signed_state::{Backend, Mirrors, ProfileStore, RepoStore};
|
||||||
use signed_ui::{Avatar, CountBadge, placeholder, status_badge};
|
use signed_ui::{Avatar, CountBadge, placeholder, status_badge};
|
||||||
use utils::{relative_time, relative_time_secs};
|
use utils::{relative_time, relative_time_secs};
|
||||||
|
|
||||||
@@ -261,7 +261,7 @@ impl PullRequestDetailView {
|
|||||||
|
|
||||||
Some(
|
Some(
|
||||||
cx.background_spawn(async move {
|
cx.background_spawn(async move {
|
||||||
let repo = ensure_repo_mirror(&addr, &clone_urls)?;
|
let repo = Mirrors::ensure(&addr, &clone_urls)?;
|
||||||
let gix_repo = repo.inner();
|
let gix_repo = repo.inner();
|
||||||
|
|
||||||
let workdir = gix_repo
|
let workdir = gix_repo
|
||||||
|
|||||||
@@ -23,9 +23,7 @@ use gpui_component::{
|
|||||||
use nostr::prelude::*;
|
use nostr::prelude::*;
|
||||||
use signed_core::{Announcement, RepoAddr};
|
use signed_core::{Announcement, RepoAddr};
|
||||||
use signed_git::{GitCache, Repo};
|
use signed_git::{GitCache, Repo};
|
||||||
use signed_state::{
|
use signed_state::{Backend, CheckoutsStore, Mirrors, RepoListStore, RepoStore};
|
||||||
Backend, CheckoutsStore, RepoListStore, RepoStore, ensure_repo_mirror, repo_mirror_path,
|
|
||||||
};
|
|
||||||
use signed_ui::{CountBadge, placeholder, ref_selector_trigger};
|
use signed_ui::{CountBadge, placeholder, ref_selector_trigger};
|
||||||
use utils::middle_truncate;
|
use utils::middle_truncate;
|
||||||
|
|
||||||
@@ -531,7 +529,7 @@ impl NewPullRequestView {
|
|||||||
let Some((base, _euc)) = self.base_repo(cx) else {
|
let Some((base, _euc)) = self.base_repo(cx) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
let mirror_path = repo_mirror_path(&base);
|
let mirror_path = Mirrors::path(&base);
|
||||||
let namespace = GitCache::fork_namespace(&announcement);
|
let namespace = GitCache::fork_namespace(&announcement);
|
||||||
let clone_urls = announcement.clone.clone();
|
let clone_urls = announcement.clone.clone();
|
||||||
|
|
||||||
@@ -566,7 +564,7 @@ impl NewPullRequestView {
|
|||||||
let clone_urls = clone_urls.clone();
|
let clone_urls = clone_urls.clone();
|
||||||
let mirror_path = mirror_path.clone();
|
let mirror_path = mirror_path.clone();
|
||||||
async move {
|
async move {
|
||||||
ensure_repo_mirror(&base, &base_clone_urls)?;
|
Mirrors::ensure(&base, &base_clone_urls)?;
|
||||||
let repo = Repo::open(&mirror_path)?;
|
let repo = Repo::open(&mirror_path)?;
|
||||||
|
|
||||||
// Prune stale imports of any fork. Then import this fork's heads under its namespace.
|
// Prune stale imports of any fork. Then import this fork's heads under its namespace.
|
||||||
|
|||||||
@@ -25,9 +25,8 @@ use nostr::prelude::{RelayUrl, ToBech32, Url};
|
|||||||
use signed_core::{Announcement, RepoAddr, RepoStatus};
|
use signed_core::{Announcement, RepoAddr, RepoStatus};
|
||||||
use signed_git::{FileCommit, GitCache, Repo};
|
use signed_git::{FileCommit, GitCache, Repo};
|
||||||
use signed_state::{
|
use signed_state::{
|
||||||
Backend, CheckoutStatus, CheckoutsStore, LocalReposStore, Nip34Binding, Nip34Kind,
|
Backend, CheckoutStatus, CheckoutsStore, LocalReposStore, Mirrors, Nip34Binding, Nip34Kind,
|
||||||
ProfileStore, RepoListStore, RepoStore, ensure_repo_mirror, open_repo_mirror,
|
ProfileStore, RepoListStore, RepoStore,
|
||||||
pr_proposes_checkout,
|
|
||||||
};
|
};
|
||||||
use signed_ui::{
|
use signed_ui::{
|
||||||
Avatar, CountBadge, DropdownButton, PixelAvatar, copy_row, menu_copy_row, ref_selector_trigger,
|
Avatar, CountBadge, DropdownButton, PixelAvatar, copy_row, menu_copy_row, ref_selector_trigger,
|
||||||
@@ -359,7 +358,7 @@ impl RepoDetailView {
|
|||||||
let disk = {
|
let disk = {
|
||||||
let addr = addr.clone();
|
let addr = addr.clone();
|
||||||
cx.background_spawn(async move {
|
cx.background_spawn(async move {
|
||||||
match open_repo_mirror(&addr)? {
|
match Mirrors::open(&addr)? {
|
||||||
Some(repo) => Ok(Some(load_repo_data(&repo)?)),
|
Some(repo) => Ok(Some(load_repo_data(&repo)?)),
|
||||||
None => Ok(None),
|
None => Ok(None),
|
||||||
}
|
}
|
||||||
@@ -376,7 +375,7 @@ impl RepoDetailView {
|
|||||||
let addr = addr.clone();
|
let addr = addr.clone();
|
||||||
let clone_urls = clone_urls.clone();
|
let clone_urls = clone_urls.clone();
|
||||||
cx.background_spawn(async move {
|
cx.background_spawn(async move {
|
||||||
let repo = ensure_repo_mirror(&addr, &clone_urls)?;
|
let repo = Mirrors::ensure(&addr, &clone_urls)?;
|
||||||
load_repo_data(&repo)
|
load_repo_data(&repo)
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
@@ -403,7 +402,7 @@ impl RepoDetailView {
|
|||||||
let addr = addr.clone();
|
let addr = addr.clone();
|
||||||
|
|
||||||
cx.background_spawn(async move {
|
cx.background_spawn(async move {
|
||||||
let Some(repo) = open_repo_mirror(&addr)? else {
|
let Some(repo) = Mirrors::open(&addr)? else {
|
||||||
return Ok::<_, Error>(None);
|
return Ok::<_, Error>(None);
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -1506,8 +1505,12 @@ impl RepoDetailView {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
for pr in &store.pull_requests {
|
for pr in &store.pull_requests {
|
||||||
if pr_proposes_checkout(pr, store.status_of(pr) == RepoStatus::Open, user, &status)
|
if CheckoutsStore::pr_proposes_checkout(
|
||||||
{
|
pr,
|
||||||
|
store.status_of(pr) == RepoStatus::Open,
|
||||||
|
user,
|
||||||
|
&status,
|
||||||
|
) {
|
||||||
continue 'status;
|
continue 'status;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ use nostr::prelude::RelayUrl;
|
|||||||
use signed_core::{Announcement, RepoAddr};
|
use signed_core::{Announcement, RepoAddr};
|
||||||
use signed_state::{
|
use signed_state::{
|
||||||
Backend, BackendEvent, CheckoutsStore, LocalReposStore, Nip34Binding, Nip34Kind, Profile,
|
Backend, BackendEvent, CheckoutsStore, LocalReposStore, Nip34Binding, Nip34Kind, Profile,
|
||||||
ProfileStore, RepoListStore, ResolvedLocalRepo, resolve_local_repos,
|
ProfileStore, RepoListStore, ResolvedLocalRepo,
|
||||||
};
|
};
|
||||||
use signed_ui::{Avatar, NavItem, PixelAvatar, title_bar_drag_handlers};
|
use signed_ui::{Avatar, NavItem, PixelAvatar, title_bar_drag_handlers};
|
||||||
|
|
||||||
@@ -136,7 +136,7 @@ impl SidebarPanel {
|
|||||||
|
|
||||||
let local = LocalReposStore::global(cx).read(cx);
|
let local = LocalReposStore::global(cx).read(cx);
|
||||||
let local_repos =
|
let local_repos =
|
||||||
resolve_local_repos(&local.repos, &repo_list.announcements, &announcements);
|
LocalReposStore::resolve(&local.repos, &repo_list.announcements, &announcements);
|
||||||
|
|
||||||
(announcements, local_repos, local.scanning)
|
(announcements, local_repos, local.scanning)
|
||||||
};
|
};
|
||||||
|
|||||||
Reference in New Issue
Block a user