Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
211 changes: 204 additions & 7 deletions crates/threadlane-git/src/git.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::{Mutex, OnceLock};
use std::sync::{Arc, Condvar, Mutex, OnceLock};
use std::time::{Duration, Instant};

use crate::error::GitError;
Expand All @@ -14,6 +14,7 @@ use crate::types::{
#[cfg(test)]
thread_local! {
pub(crate) static COMMAND_SPAWNS: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
pub(crate) static REMOTE_FETCH_RUNS: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
}

const REPOSITORY_METADATA_TTL: Duration = Duration::from_secs(60);
Expand Down Expand Up @@ -204,11 +205,209 @@ fn apply_numstats(work_dir: &Path, status: &mut GitStatus) {
}
}

pub fn sync_remote(work_dir: &Path) -> Result<(), GitError> {
/// Background remote syncs share one fetch per repository per window.
/// Every refresh path (project switches, panel opens, watcher-driven
/// reloads) funnels here, so a burst of triggers costs one remote
/// round-trip instead of one per caller.
const REMOTE_SYNC_TTL: Duration = Duration::from_secs(60);

/// Fetch bookkeeping for one repository. Linked worktrees share
/// `refs/remotes/*`, so they share a slot too.
#[derive(Default)]
struct RemoteFetchSlot {
state: Mutex<RemoteFetchState>,
changed: Condvar,
}

#[derive(Default)]
struct RemoteFetchState {
in_flight: bool,
/// Callers blocked on `changed`, observable to tests.
waiters: usize,
/// The completed run's outcome; errors reach queued callers but are
/// never freshness-stamped, so a failed fetch does not suppress retries.
last_completed: Option<(Instant, Result<(), String>)>,
}

static REMOTE_FETCH_SLOTS: OnceLock<Mutex<HashMap<(PathBuf, String), Arc<RemoteFetchSlot>>>> =
OnceLock::new();

/// A bare `git fetch` reads `branch.<name>.remote`, falling back to
/// `origin`. Sibling worktrees on branches tracking different remotes must
/// not share a slot, or one fetch satisfies callers expecting the other.
fn fetch_remote_name(work_dir: &Path) -> String {
if let Ok(branch) = command(work_dir, &["symbolic-ref", "--quiet", "--short", "HEAD"]) {
let key = format!("branch.{}.remote", branch.trim());
if let Ok(remote) = command(work_dir, &["config", "--get", &key]) {
let remote = remote.trim();
if !remote.is_empty() {
return remote.to_owned();
}
}
}
"origin".to_owned()
}

fn remote_fetch_slot(work_dir: &Path) -> Arc<RemoteFetchSlot> {
// `--git-common-dir` resolves to the shared refs dir for linked
// worktrees, so callers in sibling worktrees land on one slot.
let dir = command(work_dir, &["rev-parse", "--git-common-dir"])
.ok()
.map(|output| {
let dir = output.trim();
let path = if Path::new(dir).is_absolute() {
PathBuf::from(dir)
} else {
work_dir.join(dir)
};
path.canonicalize().unwrap_or(path)
})
.filter(|path| !path.as_os_str().is_empty())
.unwrap_or_else(|| repository_key(work_dir));
let key = (dir, fetch_remote_name(work_dir));
REMOTE_FETCH_SLOTS
.get_or_init(|| Mutex::new(HashMap::new()))
.lock()
.unwrap_or_else(|error| error.into_inner())
.entry(key)
.or_default()
.clone()
}

fn run_remote_fetch(
work_dir: &Path,
slot: &RemoteFetchSlot,
ttl: Duration,
fetch_op: impl FnOnce(&Path) -> Result<(), GitError>,
) -> Result<(), GitError> {
{
let mut state = slot.state.lock().unwrap_or_else(|e| e.into_inner());
loop {
if state.in_flight {
state.waiters += 1;
state = slot.changed.wait(state).unwrap_or_else(|e| e.into_inner());
state.waiters -= 1;
if !state.in_flight {
// The run just completed is fresh by definition: its
// result satisfies a queued caller on any window.
if let Some((_, result)) = &state.last_completed {
return result
.clone()
.map_err(|message| GitError::new(work_dir, message));
}
}
continue;
}
if let Some((fetched_at, result)) = &state.last_completed {
if result.is_ok() && fetched_at.elapsed() <= ttl {
return result
.clone()
.map_err(|message| GitError::new(work_dir, message));
}
}
state.in_flight = true;
break;
}
}
let outcome = fetch_op(work_dir);
let mut state = slot.state.lock().unwrap_or_else(|e| e.into_inner());
state.in_flight = false;
state.last_completed = Some((
Instant::now(),
outcome.clone().map_err(|error| error.message.clone()),
));
slot.changed.notify_all();
outcome
}

fn git_fetch(work_dir: &Path) -> Result<(), GitError> {
#[cfg(test)]
REMOTE_FETCH_RUNS.set(REMOTE_FETCH_RUNS.get() + 1);
command(work_dir, &["fetch", "--prune", "--quiet"])?;
Ok(())
}

/// `fetch --prune` when the repository's last successful fetch is older
/// than [`REMOTE_SYNC_TTL`]; a fresher run or one already in flight
/// answers the caller without another remote round-trip.
pub fn sync_remote(work_dir: &Path) -> Result<(), GitError> {
let slot = remote_fetch_slot(work_dir);
run_remote_fetch(work_dir, &slot, REMOTE_SYNC_TTL, git_fetch)
}

/// Explicit fetch (the panel's Fetch action): always contacts the remote
/// but still coalesces with a run already in flight instead of racing it
/// on `.git` locks.
pub fn fetch(work_dir: &Path) -> Result<(), GitError> {
let slot = remote_fetch_slot(work_dir);
run_remote_fetch(work_dir, &slot, Duration::ZERO, git_fetch)
}

/// A completed pull already refreshed remote-tracking refs; let the next
/// background sync skip its fetch.
fn note_remote_fetched(work_dir: &Path) {
let slot = remote_fetch_slot(work_dir);
slot.state
.lock()
.unwrap_or_else(|e| e.into_inner())
.last_completed = Some((Instant::now(), Ok(())));
}

#[cfg(test)]
mod fetch_tests {
use super::{run_remote_fetch, GitError, RemoteFetchSlot};
use std::path::Path;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Barrier};
use std::time::Duration;

#[test]
fn queued_callers_share_the_in_flight_run() {
let slot = RemoteFetchSlot::default();
let runs = Arc::new(AtomicUsize::new(0));
let entered = Arc::new(Barrier::new(2));
let release = Arc::new(Barrier::new(2));
std::thread::scope(|scope| {
let leader = {
let runs = runs.clone();
let entered = entered.clone();
let release = release.clone();
let slot = &slot;
scope.spawn(move || {
run_remote_fetch(Path::new("/repo"), slot, Duration::ZERO, move |_| {
runs.fetch_add(1, Ordering::SeqCst);
entered.wait();
release.wait();
Ok(())
})
})
};
entered.wait();
let waiters: Vec<_> = (0..4)
.map(|_| {
let slot = &slot;
scope.spawn(move || {
run_remote_fetch(Path::new("/repo"), slot, Duration::ZERO, |_| {
Err(GitError::new("/repo", "must not run"))
})
})
})
.collect();
// Release the leader only once every waiter is blocked on the
// condvar, so a descheduled waiter can't restart a fresh run.
while slot.state.lock().unwrap().waiters < waiters.len() {
std::thread::yield_now();
}
release.wait();
for waiter in waiters {
waiter.join().unwrap().unwrap();
}
leader.join().unwrap().unwrap();
});
assert_eq!(runs.load(Ordering::SeqCst), 1);
}
}

fn repository_metadata(work_dir: &Path) -> RepositoryMetadata {
let key = repository_key(work_dir);
let now = Instant::now();
Expand All @@ -233,10 +432,6 @@ fn repository_metadata(work_dir: &Path) -> RepositoryMetadata {
metadata
}

pub fn fetch(work_dir: &Path) -> Result<(), GitError> {
sync_remote(work_dir)
}

pub(crate) fn list_branches_detailed(
work_dir: &Path,
provided_default_branch: Option<&str>,
Expand Down Expand Up @@ -1076,7 +1271,9 @@ pub fn push(work_dir: &Path) -> Result<(), GitError> {
}

pub fn pull(work_dir: &Path) -> Result<String, GitError> {
command(work_dir, &["pull", "--ff-only"])
let output = command(work_dir, &["pull", "--ff-only"])?;
note_remote_fetched(work_dir);
Ok(output)
}

pub fn merge(work_dir: &Path, branch: &str) -> Result<String, GitError> {
Expand Down
100 changes: 100 additions & 0 deletions crates/threadlane-git/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1697,6 +1697,106 @@ fn worktree_bases_refreshes_refs_and_selects_only_existing_defaults() {
assert_eq!(worktree_bases(local.path()).unwrap(), (default, branches));
}

fn init_fetch_repo_pair() -> (tempfile::TempDir, tempfile::TempDir) {
let remote = tempdir().unwrap();
let local = tempdir().unwrap();
for dir in [remote.path(), local.path()] {
run_git(dir, &["init", "-q", "-b", "main"]);
run_git(dir, &["config", "user.email", "test@example.com"]);
run_git(dir, &["config", "user.name", "Test"]);
}
run_git(
remote.path(),
&["commit", "--allow-empty", "-qm", "initial"],
);
run_git(
local.path(),
&["remote", "add", "origin", remote.path().to_str().unwrap()],
);
(remote, local)
}

#[test]
fn sync_remote_coalesces_fetches_within_a_window() {
let (_remote, local) = init_fetch_repo_pair();
REMOTE_FETCH_RUNS.set(0);
sync_remote(local.path()).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 1);
// Within the window the remote state is shared, not refetched.
sync_remote(local.path()).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 1);
// An explicit Fetch always runs and re-arms the window.
fetch(local.path()).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 2);
sync_remote(local.path()).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 2);
}

#[test]
fn sync_remote_shares_one_fetch_across_linked_worktrees() {
let (_remote, local) = init_fetch_repo_pair();
run_git(local.path(), &["commit", "--allow-empty", "-qm", "local"]);
let sibling = local.path().join("sibling");
run_git(
local.path(),
&[
"worktree",
"add",
"-q",
sibling.to_str().unwrap(),
"-b",
"sibling",
],
);
REMOTE_FETCH_RUNS.set(0);
sync_remote(local.path()).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 1);
// A linked worktree shares refs/remotes with the root, so its sync
// is already fresh — no second remote call.
sync_remote(&sibling).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 1);
}

#[test]
fn sync_remote_retries_after_a_failed_fetch() {
let (_remote, local) = init_fetch_repo_pair();
run_git(
local.path(),
&[
"remote",
"set-url",
"origin",
"/nonexistent/threadlane-test-remote",
],
);
REMOTE_FETCH_RUNS.set(0);
assert!(sync_remote(local.path()).is_err());
// A failure is not freshness-stamped: the next sync retries.
assert!(sync_remote(local.path()).is_err());
assert_eq!(REMOTE_FETCH_RUNS.get(), 2);
}

#[test]
fn pull_marks_remote_refs_fresh_for_background_syncs() {
let (remote, _keep) = init_fetch_repo_pair();
let parent = tempdir().unwrap();
run_git(
parent.path(),
&[
"clone",
"-q",
remote.path().to_str().unwrap(),
"local",
],
);
let local = parent.path().join("local");
pull(&local).unwrap();
REMOTE_FETCH_RUNS.set(0);
// The pull already refreshed remote-tracking refs.
sync_remote(&local).unwrap();
assert_eq!(REMOTE_FETCH_RUNS.get(), 0);
}

#[test]
fn list_project_files_reports_tracked_and_nonignored_untracked() {
let dir = tempdir().unwrap();
Expand Down
Loading