feat: add event fetching strategy #24

Merged
reya merged 10 commits from optimize into master 2026-09-27 02:50:12 +00:00
4 changed files with 15 additions and 73 deletions
Showing only changes of commit 0ee67d93b6 - Show all commits
+2
View File
@@ -81,6 +81,8 @@ Both `cx.spawn` and `cx.background_spawn` return a `Task<R>`, which is a future
A task which doesn't do anything but provide a value can be created with `Task::ready(value)`. A task which doesn't do anything but provide a value can be created with `Task::ready(value)`.
Prefer keeping a task in a field over `.detach()` when the work belongs to a view or store. A detached task outlives the view that started it (for example, a repository panel the user has closed), while a task stored in a field is cancelled when that struct drops. A finished `Task` held in a container is not reaped by GPUI - it stays alive until its handle is dropped - so if the container can accumulate many runs, either drop the finished handles before pushing a new one, or hold a single `Option<Task<_>>` and replace it.
## Elements ## Elements
The `Render` trait is used to render some state into an element tree that is laid out using flexbox layout. An `Entity<T>` where `T` implements `Render` is sometimes called a "view". The `Render` trait is used to render some state into an element tree that is laid out using flexbox layout. An `Entity<T>` where `T` implements `Render` is sometimes called a "view".
+4 -46
View File
@@ -17,15 +17,9 @@ use crate::repos::RepoListStore;
const REFRESH_DEBOUNCE: Duration = Duration::from_millis(300); const REFRESH_DEBOUNCE: Duration = Duration::from_millis(300);
/// How often the statuses are recomputed against the local refs. /// How often the statuses are recomputed against the local refs.
///
/// A commit lands in a checkout long before the remote reconciliation cadence,
/// so this fast pass surfaces ready-to-push and ready-to-contribute checkouts
/// within a second or two. It reads the tracking refs only, no network.
const LOCAL_POLL: Duration = Duration::from_secs(2); const LOCAL_POLL: Duration = Duration::from_secs(2);
/// How often a full pass refreshes the remotes while any repository panel is open. /// How often a full pass refreshes the remotes while any repository panel is open.
const STATUS_POLL: Duration = Duration::from_secs(15); const STATUS_POLL: Duration = Duration::from_secs(15);
/// Remote refresh interval for the `ready to push` badges of the user's own repositories. /// Remote refresh interval for the `ready to push` badges of the user's own repositories.
const PUSH_POLL: Duration = Duration::from_secs(60); const PUSH_POLL: Duration = Duration::from_secs(60);
@@ -46,10 +40,6 @@ pub struct CheckoutStatus {
/// Commit the branch points at, for tip-based PR dedupe. /// Commit the branch points at, for tip-based PR dedupe.
pub head: String, pub head: String,
/// What the branch is compared against. /// What the branch is compared against.
/// For ready-to-contribute statuses, the announced HEAD branch.
/// The fallbacks are `main`, then the first local branch.
/// For ready-to-push statuses, the remote-tracking ref.
/// Unpushed commits are counted against it.
/// ///
/// It is `refs/remotes/origin/<branch>`, else `origin/HEAD` for new branches. /// It is `refs/remotes/origin/<branch>`, else `origin/HEAD` for new branches.
pub base: String, pub base: String,
@@ -67,11 +57,6 @@ struct Remembered {
} }
/// Global store of local-checkout associations and per-checkout statuses. /// Global store of local-checkout associations and per-checkout statuses.
///
/// Readers (the sidebar rows, the repository panels) observe this store and
/// derive what they display from their own snapshots, so publishing needs no
/// fine-grained entities: the store notifies when a slice changed and each
/// reader re-derives only what it shows.
pub struct CheckoutsStore { pub struct CheckoutsStore {
/// Checkout paths per announced repository. /// Checkout paths per announced repository.
by_repo: HashMap<RepoAddr, Vec<PathBuf>>, by_repo: HashMap<RepoAddr, Vec<PathBuf>>,
@@ -144,6 +129,7 @@ impl CheckoutsStore {
this.statuses.clear(); this.statuses.clear();
this.push_statuses.clear(); this.push_statuses.clear();
cx.notify(); cx.notify();
this.refresh(cx); this.refresh(cx);
} }
})); }));
@@ -282,9 +268,6 @@ impl CheckoutsStore {
} }
/// Re-resolve the associations and the requested statuses. /// Re-resolve the associations and the requested statuses.
///
/// Requests arriving while a pass runs fold into a follow-up, requests
/// arriving while the debounce timer is pending are dropped.
pub fn refresh(&mut self, cx: &mut Context<Self>) { pub fn refresh(&mut self, cx: &mut Context<Self>) {
if self.debounce_pending || self.refresh.request() != RefreshRequest::Schedule { if self.debounce_pending || self.refresh.request() != RefreshRequest::Schedule {
return; return;
@@ -300,12 +283,6 @@ impl CheckoutsStore {
} }
/// One full resolve and apply cycle, the debounced entry point. /// One full resolve and apply cycle, the debounced entry point.
///
/// Re-resolves the associations from the settings, the scan and the
/// announcements, then recomputes the requested statuses against freshly
/// fetched remotes. Full passes run on every input change and on the
/// remote reconciliation cadence ([`Self::local_tick`]); they also restart
/// the fast local pass.
fn run_refresh(&mut self, cx: &mut Context<Self>) { fn run_refresh(&mut self, cx: &mut Context<Self>) {
self.debounce_pending = false; self.debounce_pending = false;
self.refresh.begin(); self.refresh.begin();
@@ -431,10 +408,6 @@ impl CheckoutsStore {
} }
/// Schedule the fast local status pass, unless one is already pending. /// Schedule the fast local status pass, unless one is already pending.
///
/// Every [`LOCAL_POLL`] the pass recomputes the requested statuses against
/// the local refs, with no network, so a new commit in a checkout surfaces in
/// a second or two instead of at the next remote reconciliation.
fn schedule_local_pass(&mut self, cx: &mut Context<Self>) { fn schedule_local_pass(&mut self, cx: &mut Context<Self>) {
if self.local_pending { if self.local_pending {
return; return;
@@ -452,10 +425,6 @@ impl CheckoutsStore {
} }
/// The fast local status pass. /// The fast local status pass.
///
/// Recomputes the statuses against the local refs; when the remote
/// reconciliation cadence elapsed, it runs a full pass instead so pushes
/// made elsewhere do not linger as `to push`.
fn local_tick(&mut self, cx: &mut Context<Self>) { fn local_tick(&mut self, cx: &mut Context<Self>) {
// Nothing watched: the chain idles out until a new request restarts it. // Nothing watched: the chain idles out until a new request restarts it.
if self.status_requested.is_empty() && self.push_requested.is_empty() { if self.status_requested.is_empty() && self.push_requested.is_empty() {
@@ -490,10 +459,6 @@ impl CheckoutsStore {
} }
/// Recompute the requested statuses against the tracking refs only. /// Recompute the requested statuses against the tracking refs only.
///
/// The refs were last refreshed by a full pass. Comparing against them is
/// enough to pick up new local commits, and skipping the network keeps
/// this pass cheap enough to run every [`LOCAL_POLL`].
fn run_local_statuses(&mut self, cx: &mut Context<Self>) { fn run_local_statuses(&mut self, cx: &mut Context<Self>) {
let associations = self.by_repo.clone(); let associations = self.by_repo.clone();
@@ -635,10 +600,6 @@ fn checkout_status(path: &Path, announced_head: Option<&str>) -> Option<Checkout
} }
/// 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.
///
/// `fetch` refreshes the remote heads first, so a full pass sees pushes made
/// elsewhere; the fast local pass skips it and compares against the tracking
/// refs left by the last full pass, which is enough to detect local commits.
fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> { fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
if signed_git::worktree_dirty(path) { if signed_git::worktree_dirty(path) {
return None; return None;
@@ -649,8 +610,6 @@ fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
let origin = signed_git::origin_url(path).ok().flatten()?; let origin = signed_git::origin_url(path).ok().flatten()?;
if fetch { if fetch {
// Refresh the remote heads first.
// Commits made elsewhere or pushed from another machine must not linger as `to push`.
signed_git::fetch_repo_refs(path, &[origin], "+refs/heads/*:refs/remotes/origin/*").ok(); signed_git::fetch_repo_refs(path, &[origin], "+refs/heads/*:refs/remotes/origin/*").ok();
} }
@@ -678,10 +637,6 @@ fn checkout_push_status(path: &Path, fetch: bool) -> Option<CheckoutStatus> {
} }
/// Compute the requested statuses against the checkout paths of `associations`. /// Compute the requested statuses against the checkout paths of `associations`.
///
/// Shared by the full and the local pass. `fetch` refreshes the checkouts'
/// remote heads first, so the full pass sees remote moves; the fast local
/// pass reads the tracking refs only, which is enough to detect local commits.
fn compute_statuses( fn compute_statuses(
associations: &HashMap<RepoAddr, Vec<PathBuf>>, associations: &HashMap<RepoAddr, Vec<PathBuf>>,
requested: &[(RepoAddr, Option<String>)], requested: &[(RepoAddr, Option<String>)],
@@ -737,12 +692,14 @@ pub fn pr_proposes_checkout(
if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user { if pr.kind != Kind::GitPullRequest || !open || pr.pubkey != user {
return false; return false;
} }
let branch_matches = pr let branch_matches = pr
.tags .tags
.iter() .iter()
.find(|t| t.kind() == "branch-name") .find(|t| t.kind() == "branch-name")
.and_then(|t| t.content()) .and_then(|t| t.content())
.is_some_and(|name| name == checkout.branch); .is_some_and(|name| name == checkout.branch);
// A renamed branch falls back to the proposed tip commit. // A renamed branch falls back to the proposed tip commit.
let tip_matches = pr let tip_matches = pr
.tags .tags
@@ -750,6 +707,7 @@ pub fn pr_proposes_checkout(
.find(|t| t.kind() == "c") .find(|t| t.kind() == "c")
.and_then(|t| t.content()) .and_then(|t| t.content())
.is_some_and(|tip| tip == checkout.head); .is_some_and(|tip| tip == checkout.head);
branch_matches || tip_matches branch_matches || tip_matches
} }
-4
View File
@@ -18,10 +18,6 @@ impl RefreshGate {
self.running self.running
} }
/// A new refresh request arrived.
///
/// Folded into a follow-up run while one is in flight, otherwise the
/// caller starts the run itself.
pub fn request(&mut self) -> RefreshRequest { pub fn request(&mut self) -> RefreshRequest {
if self.running { if self.running {
self.dirty = true; self.dirty = true;
+9 -23
View File
@@ -20,7 +20,6 @@ use crate::backend::{
}; };
use crate::checkouts::CheckoutsStore; use crate::checkouts::CheckoutsStore;
use crate::git_store::ensure_repo_mirror; use crate::git_store::ensure_repo_mirror;
use crate::refresh::{RefreshGate, RefreshRequest};
use crate::repos::RepoListStore; use crate::repos::RepoListStore;
/// Maximum size of one patch event. /// Maximum size of one patch event.
@@ -31,7 +30,6 @@ const MAX_PATCH_EVENT_BYTES: usize = 60 * 1024;
/// Per-repository store. /// Per-repository store.
/// ///
/// Holds the announcement, state, issues, patches, PRs, comments and resolved statuses. /// Holds the announcement, state, issues, patches, PRs, comments and resolved statuses.
/// Always derived from the local database.
pub struct RepoStore { pub struct RepoStore {
/// NIP-34 address. `None` while the repository is local-only. /// NIP-34 address. `None` while the repository is local-only.
addr: Option<RepoAddr>, addr: Option<RepoAddr>,
@@ -89,7 +87,8 @@ pub struct RepoStore {
/// ///
/// Avoids re-running the maintainer Auto sync on every refresh. /// Avoids re-running the maintainer Auto sync on every refresh.
synced_maintainers: HashSet<PublicKey>, synced_maintainers: HashSet<PublicKey>,
refresh: RefreshGate, /// In-flight refresh tasks.
tasks: Vec<Task<Result<(), Error>>>,
/// Backend subscription of an announced repository. `None` while local-only. /// Backend subscription of an announced repository. `None` while local-only.
_subscription: Option<Subscription>, _subscription: Option<Subscription>,
} }
@@ -138,7 +137,7 @@ impl RepoStore {
repo_relays: HashSet::new(), repo_relays: HashSet::new(),
root_fetches: HashSet::new(), root_fetches: HashSet::new(),
synced_maintainers: HashSet::new(), synced_maintainers: HashSet::new(),
refresh: RefreshGate::default(), tasks: Vec::new(),
_subscription: Some(subscription), _subscription: Some(subscription),
} }
} }
@@ -167,7 +166,7 @@ impl RepoStore {
repo_relays: HashSet::new(), repo_relays: HashSet::new(),
root_fetches: HashSet::new(), root_fetches: HashSet::new(),
synced_maintainers: HashSet::new(), synced_maintainers: HashSet::new(),
refresh: RefreshGate::default(), tasks: Vec::new(),
_subscription: None, _subscription: None,
} }
} }
@@ -365,11 +364,6 @@ impl RepoStore {
if self.addr.is_none() { if self.addr.is_none() {
return; return;
} }
if self.refresh.request() != RefreshRequest::Schedule {
return;
}
self.run_refresh(cx); self.run_refresh(cx);
} }
@@ -378,8 +372,6 @@ impl RepoStore {
return; return;
}; };
self.refresh.begin();
let backend = Backend::global(cx); let backend = Backend::global(cx);
let client = backend.read(cx).client(); let client = backend.read(cx).client();
@@ -498,7 +490,7 @@ impl RepoStore {
)) ))
}); });
cx.spawn(async move |this, cx| { let task = cx.spawn(async move |this, cx| {
let ( let (
announcement, announcement,
state, state,
@@ -513,14 +505,13 @@ impl RepoStore {
Ok(data) => data, Ok(data) => data,
Err(e) => { Err(e) => {
return this.update(cx, |this, cx| { return this.update(cx, |this, cx| {
this.refresh.abort();
this.last_error = Some(e.to_string()); this.last_error = Some(e.to_string());
cx.notify(); cx.notify();
}); });
} }
}; };
let again = this.update(cx, |this, cx| { this.update(cx, |this, cx| {
let keep_hint = announcement.is_none() && !this.loaded; let keep_hint = announcement.is_none() && !this.loaded;
let first_pass = !this.loaded; let first_pass = !this.loaded;
@@ -604,17 +595,12 @@ impl RepoStore {
if changed { if changed {
cx.notify(); cx.notify();
} }
this.refresh.finish()
})?; })?;
if again {
this.update(cx, |this, cx| this.refresh(cx))?;
}
Ok(()) Ok(())
}) });
.detach();
self.tasks.push(task);
} }
/// Resolve the status of a root event, an issue, patch or PR, per NIP-34. /// Resolve the status of a root event, an issue, patch or PR, per NIP-34.