use std::collections::{HashMap, HashSet};
use std::hash::{Hash, Hasher};
use std::path::Path;
use std::process::Command;
use std::time::{Duration, Instant, UNIX_EPOCH};
use ignore::WalkBuilder;
use crate::core::RepoIdentity;
use crate::lang;
use crate::store::Store;
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub(crate) struct Stats {
pub files_seen: usize,
pub files_indexed: usize,
pub symbols: usize,
}
#[cfg(test)]
pub(crate) fn index_path(
store: &mut Store,
root: &Path,
) -> Result<Stats, Box<dyn std::error::Error>> {
index_under(store, root, &[])
}
pub(crate) fn index_under(
store: &mut Store,
root: &Path,
subdirs: &[String],
) -> Result<Stats, Box<dyn std::error::Error>> {
run_index(store, root, &[], subdirs, None, None, None)
}
fn alnum_lower(s: &str) -> String {
s.chars()
.filter(|c| c.is_alphanumeric())
.map(|c| c.to_ascii_lowercase())
.collect()
}
fn prioritize_by_path(
paths: Vec<std::path::PathBuf>,
_root: &Path,
query: Option<&str>,
) -> Vec<std::path::PathBuf> {
let needle = alnum_lower(query.unwrap_or(""));
let k = needle.len().min(4);
if k == 0 {
return paths;
}
let kgrams: std::collections::HashSet<&[u8]> = needle.as_bytes().windows(k).collect();
let mut prio = Vec::new();
let mut rest = Vec::new();
let mut stem = String::new();
for p in paths {
stem.clear();
if let Some(s) = p.file_stem() {
stem.extend(
s.to_string_lossy()
.chars()
.filter(|c| c.is_alphanumeric())
.map(|c| c.to_ascii_lowercase()),
);
}
if stem.as_bytes().windows(k).any(|w| kgrams.contains(w)) {
prio.push(p);
} else {
rest.push(p);
}
}
prio.extend(rest);
prio
}
pub(crate) fn index_budgeted(
store: &mut Store,
root: &Path,
active: &[String],
budget: Duration,
query: Option<&str>,
) -> Result<Stats, Box<dyn std::error::Error>> {
run_index(store, root, active, &[], Some(budget), query, None)
}
pub(crate) fn index_budgeted_cancellable(
store: &mut Store,
root: &Path,
active: &[String],
budget: Duration,
query: Option<&str>,
cancel: &std::sync::atomic::AtomicBool,
) -> Result<Stats, Box<dyn std::error::Error>> {
run_index(store, root, active, &[], Some(budget), query, Some(cancel))
}
const COLLECT_CAP: usize = 50_000;
fn collect_cap() -> usize {
std::env::var("RQ_COLLECT_CAP")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(COLLECT_CAP)
}
static PARSE_JOBS: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
pub(crate) fn set_parse_jobs(n: usize) {
PARSE_JOBS.store(n, std::sync::atomic::Ordering::Relaxed);
}
pub(crate) fn parse_jobs() -> usize {
let configured = PARSE_JOBS.load(std::sync::atomic::Ordering::Relaxed);
if configured > 0 {
return configured;
}
if let Some(n) = std::env::var("RQ_JOBS").ok().and_then(|v| v.parse().ok())
&& n > 0
{
return n;
}
std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
}
static PARSE_US: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
const WRITE_BATCH: usize = 512;
const WRITE_INTERVAL: Duration = Duration::from_millis(50);
struct BatchWriter<'a> {
store: &'a mut Store,
repo_id: i64,
buf: Vec<crate::store::FileSymbols>,
files: usize,
symbols: usize,
write_time: Duration,
batches: usize,
last_flush: Option<Instant>,
}
impl<'a> BatchWriter<'a> {
fn new(store: &'a mut Store, repo_id: i64) -> Self {
Self {
store,
repo_id,
buf: Vec::new(),
files: 0,
symbols: 0,
write_time: Duration::ZERO,
batches: 0,
last_flush: None,
}
}
fn push(&mut self, fs: crate::store::FileSymbols) -> Result<(), Box<dyn std::error::Error>> {
self.buf.push(fs);
if self.buf.len() >= WRITE_BATCH
|| self
.last_flush
.is_none_or(|t| t.elapsed() >= WRITE_INTERVAL)
{
self.flush()?;
}
Ok(())
}
fn flush(&mut self) -> Result<(), Box<dyn std::error::Error>> {
if !self.buf.is_empty() {
let t = Instant::now();
let (f, sy) = self.store.replace_files(self.repo_id, &self.buf)?;
self.write_time += t.elapsed();
self.batches += 1;
self.files += f;
self.symbols += sy;
self.buf.clear();
self.last_flush = Some(Instant::now());
}
Ok(())
}
}
fn git_source_candidates(root: &Path) -> Option<Vec<std::path::PathBuf>> {
if !is_git_repo(root) {
return None;
}
let globs: Vec<String> = lang::registry()
.iter()
.flat_map(|p| p.extensions().iter().map(|e| format!("*.{e}")))
.collect();
let mut cmd = Command::new("git");
cmd.arg("-C")
.arg(root)
.args(["ls-files", "-z", "--cached", "--"])
.args(&globs);
let out = cmd.output().ok()?;
if !out.status.success() {
return None;
}
Some(
out.stdout
.split(|&b| b == 0)
.filter(|s| !s.is_empty())
.map(|s| root.join(String::from_utf8_lossy(s).as_ref()))
.collect(),
)
}
fn fs_walk_candidates(
roots: Vec<std::path::PathBuf>,
deadline: Option<Instant>,
) -> impl Iterator<Item = std::path::PathBuf> {
roots.into_iter().flat_map(move |root| {
WalkBuilder::new(&root)
.filter_entry(move |_| !past(deadline))
.build()
.filter_map(Result::ok)
.filter(|e| e.file_type().is_some_and(|t| t.is_file()))
.map(ignore::DirEntry::into_path)
})
}
#[allow(clippy::too_many_arguments)]
fn stream_walk(
root: &Path,
candidates: impl Iterator<Item = std::path::PathBuf> + Send,
deadline: Option<Instant>,
cap: Option<usize>,
needle: Option<&[u8]>,
seen: HashSet<String>,
keep: impl Fn(&str, &Path) -> bool + Send,
cancel: Option<&std::sync::atomic::AtomicBool>,
mut sink: impl FnMut(crate::store::FileSymbols) -> Result<(), Box<dyn std::error::Error>>,
) -> Result<(HashSet<String>, bool), Box<dyn std::error::Error>> {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
let workers = parse_jobs();
let parse_incomplete = AtomicBool::new(false);
let (path_tx, path_rx) = std::sync::mpsc::sync_channel::<std::path::PathBuf>(1024);
let (res_tx, res_rx) = std::sync::mpsc::sync_channel::<crate::store::FileSymbols>(1024);
let path_rx = Arc::new(Mutex::new(path_rx));
let (seen, walk_finished) = std::thread::scope(|s| -> Result<_, Box<dyn std::error::Error>> {
let walk = s.spawn(move || {
let mut seen = seen;
let mut finished = true;
let mut processed = 0usize;
for path in candidates {
if past(deadline) || cancel.is_some_and(|c| c.load(Ordering::Relaxed)) {
finished = false;
break;
}
let rel = path
.strip_prefix(root)
.unwrap_or(&path)
.to_string_lossy()
.into_owned();
if !is_source(&rel) {
continue;
}
if !seen.insert(rel.clone()) {
continue; }
if !keep(&rel, &path) {
continue; }
if path_tx.send(path).is_err() {
finished = false; break;
}
processed += 1;
if cap.is_some_and(|c| processed >= c) {
finished = false;
break;
}
}
if past(deadline) {
finished = false;
}
drop(path_tx); (seen, finished)
});
let parse_incomplete = &parse_incomplete;
for _ in 0..workers {
let path_rx = Arc::clone(&path_rx);
let res_tx = res_tx.clone();
s.spawn(move || {
loop {
let got = { path_rx.lock().unwrap().recv() };
let Ok(path) = got else { break }; if past(deadline) || cancel.is_some_and(|c| c.load(Ordering::Relaxed)) {
parse_incomplete.store(true, Ordering::Relaxed); break;
}
if let Some(fs) = parse_file(root, &path, needle)
&& res_tx.send(fs).is_err()
{
break;
}
}
});
}
drop(res_tx); drop(path_rx);
let res_rx = res_rx;
while let Ok(fs) = res_rx.recv() {
sink(fs)?;
}
Ok(walk.join().unwrap())
})?;
Ok((
seen,
walk_finished && !parse_incomplete.load(Ordering::Relaxed),
))
}
fn sweep_outcome(
completed: bool,
whole_repo: bool,
seen_empty: bool,
had_stored: bool,
budgeted: bool,
) -> (bool, &'static str) {
if !whole_repo {
return (false, "warming");
}
if budgeted && completed && seen_empty && had_stored {
return (false, "warming"); }
if completed {
(true, "complete")
} else {
(false, "warming")
}
}
fn run_index(
store: &mut Store,
root: &Path,
active: &[String],
subdirs: &[String],
budget: Option<Duration>,
query: Option<&str>,
cancel: Option<&std::sync::atomic::AtomicBool>,
) -> Result<Stats, Box<dyn std::error::Error>> {
let profiling = crate::profile::enabled();
PARSE_US.store(0, std::sync::atomic::Ordering::Relaxed);
let setup_span = crate::profile::span("index: setup");
let root_display = root.canonicalize().unwrap_or_else(|_| root.to_path_buf());
let identity = budget
.and_then(|_| {
store
.identity_for_root(&root_display.to_string_lossy())
.ok()
.flatten()
})
.unwrap_or_else(|| detect_identity(root).to_string());
let branch = head_branch(root);
let repo_id = store.upsert_repository(&identity, branch.as_deref())?;
store.upsert_checkout(repo_id, &root_display.to_string_lossy(), branch.as_deref())?;
for stale in store.checkout_roots(repo_id).unwrap_or_default() {
if !Path::new(&stale).exists() {
let _ = store.forget_checkout(&stale);
}
}
let stored = store.file_mtimes(repo_id)?;
let coverage_mark = store.coverage_mark(repo_id)?;
drop(setup_span);
let mut seen: HashSet<String> = HashSet::new();
if stored.is_empty() {
store.suspend_name_index(repo_id)?;
}
let mut active_to_parse: Vec<std::path::PathBuf> = Vec::new();
for rel in active {
note_candidate(
root,
&root.join(rel),
&stored,
&mut seen,
&mut active_to_parse,
);
}
let mut active_span = crate::profile::span("index: active files");
let (active_parsed, _) = parse_files(root, &active_to_parse, None, None);
let (mut files_indexed, mut symbols) = store.replace_files(repo_id, &active_parsed)?;
active_span.note(|| format!("{} file(s)", active_parsed.len()));
drop(active_span);
let walk_roots: Vec<std::path::PathBuf> = if subdirs.is_empty() {
vec![root.to_path_buf()]
} else {
subdirs.iter().map(|s| root.join(s)).collect()
};
let mut enum_span = crate::profile::span("index: enumerate");
let git_candidates = budget
.and_then(|_| git_source_candidates(root))
.filter(|paths| !paths.is_empty());
enum_span.note(|| match &git_candidates {
Some(p) => format!("git ls-files, {} path(s)", p.len()),
None => "filesystem walk (lazy — time lands in walk+parse+write)".to_string(),
});
drop(enum_span);
let git_candidates = git_candidates.map(|paths| prioritize_by_path(paths, root, query));
let deadline = budget.map(|b| Instant::now() + b);
let cap = budget.map(|_| collect_cap());
let stored_ref = &stored;
let changed = move |rel: &str, path: &Path| match stored_ref.get(rel) {
Some(&Some(m)) => Some(m) != file_mtime(path),
_ => true, };
let skipped = std::sync::atomic::AtomicU64::new(0);
let stream_start = Instant::now();
let mut fused_span = crate::profile::span("index: walk+parse+write");
let (seen, completed, walk_files, walk_symbols, write_time, batches) = {
let mut writer = BatchWriter::new(&mut *store, repo_id);
let mut demanded: HashSet<String> = HashSet::new();
let needle = query.and_then(crate::search::literal_leaf);
if let (Some(paths), Some(needle)) = (&git_candidates, needle) {
let mut demand_span = crate::profile::span("index: demand scan");
stream_walk(
root,
paths.iter().cloned(),
deadline,
None,
Some(needle.as_bytes()),
seen.clone(),
changed,
cancel,
|fs| {
demanded.insert(fs.path.clone());
writer.push(fs)
},
)?;
writer.flush()?;
demand_span.note(|| format!("{} file(s) contain {needle:?}", demanded.len()));
}
let candidates: Box<dyn Iterator<Item = std::path::PathBuf> + Send> = match git_candidates {
Some(paths) => Box::new(paths.into_iter()),
None => Box::new(fs_walk_candidates(walk_roots, deadline)),
};
let demanded = &demanded;
let skipped = &skipped;
let keep = move |rel: &str, path: &Path| {
if demanded.contains(rel) {
return false;
}
let changed = changed(rel, path);
if profiling && !changed {
skipped.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
changed
};
let (seen, completed) = stream_walk(
root,
candidates,
deadline,
cap,
None,
seen,
keep,
cancel,
|fs| writer.push(fs),
)?;
writer.flush()?;
(
seen,
completed,
writer.files,
writer.symbols,
writer.write_time,
writer.batches,
)
};
fused_span.note(|| format!("{walk_files} file(s), {walk_symbols} symbol(s)"));
drop(fused_span);
crate::profile::record(
"index: parse (worker cpu)",
Duration::from_micros(PARSE_US.load(std::sync::atomic::Ordering::Relaxed)),
|| format!("summed across {} worker(s)", parse_jobs()),
);
crate::profile::record("index: store writes", write_time, || {
format!("{batches} batch(es), serialized")
});
if crate::trace::enabled() {
let elapsed = stream_start.elapsed();
crate::trace!(
"walk+parse+write {} file(s)/{} symbol(s) in {} ms ({} ms in store writes, {} parse jobs)",
walk_files,
walk_symbols,
elapsed.as_millis(),
write_time.as_millis(),
parse_jobs(),
);
}
{
let mut span = crate::profile::span("index: name index");
let rebuilt = store.maintain_name_index(repo_id)?;
span.note(|| if rebuilt { "rebuilt" } else { "current" }.to_string());
}
files_indexed += walk_files;
symbols += walk_symbols;
let stats = Stats {
files_seen: seen.len(),
files_indexed,
symbols,
};
if profiling {
crate::profile::count("files seen", stats.files_seen as u64);
crate::profile::count("files parsed", stats.files_indexed as u64);
crate::profile::count(
"files skipped (mtime)",
skipped.load(std::sync::atomic::Ordering::Relaxed),
);
crate::profile::count("symbols", stats.symbols as u64);
crate::profile::count("batches", batches as u64);
crate::profile::count("parse jobs", parse_jobs() as u64);
}
let whole_repo = subdirs.is_empty();
let (finalize, status) = sweep_outcome(
completed,
whole_repo,
seen.is_empty(),
!stored.is_empty(),
budget.is_some(),
);
if finalize {
let mut reconcile_span = crate::profile::span("index: reconcile");
let mut forgotten = 0;
for path in stored.keys() {
if !seen.contains(path) {
store.forget_file(repo_id, path)?;
forgotten += 1;
}
}
reconcile_span.note(|| format!("{forgotten} file(s) forgotten"));
drop(reconcile_span);
if forgotten > 0 {
crate::trace!(
"reconcile {}: forgot {forgotten} file(s) not seen on disk",
crate::trace::abbrev(&root_display)
);
}
if let Some(head) = head_state(root) {
let _ = store.set_indexed_head(repo_id, &head);
if stats.files_indexed > 0 {
let edited: Vec<String> = dirty_files(root)
.into_iter()
.filter(|f| is_source(f))
.collect();
let _ = store.set_edited_files(repo_id, &edited);
}
}
}
if stats.files_indexed > 0 && repo_root(root).is_some_and(|r| r == root_display) {
let _span = crate::profile::span("index: git metadata");
capture_commit_times(store, repo_id, root);
}
let status = if status == "complete" && !store.repo_has_files(repo_id).unwrap_or(false) {
"warming"
} else {
status
};
let recorded = store.set_coverage_since(
repo_id,
stats.files_seen as i64,
stats.files_indexed as i64,
status,
&coverage_mark,
)?;
if !recorded {
crate::trace!(
"coverage {}: kept the `complete` another pass recorded during this one",
crate::trace::abbrev(&root_display)
);
}
crate::trace!(
"index {} (budget {budget:?}): {} seen, {} indexed, {} symbols → {status}",
crate::trace::abbrev(&root_display),
stats.files_seen,
stats.files_indexed,
stats.symbols,
);
Ok(stats)
}
fn note_candidate(
root: &Path,
file: &Path,
stored: &HashMap<String, Option<i64>>,
seen: &mut HashSet<String>,
to_parse: &mut Vec<std::path::PathBuf>,
) {
let rel = file
.strip_prefix(root)
.unwrap_or(file)
.to_string_lossy()
.into_owned();
if !is_source(&rel) {
return;
}
if !seen.insert(rel.clone()) {
return; }
if let Some(&Some(m)) = stored.get(&rel)
&& Some(m) == file_mtime(file)
{
return;
}
to_parse.push(file.to_path_buf());
}
fn parse_file(
root: &Path,
file: &Path,
needle: Option<&[u8]>,
) -> Option<crate::store::FileSymbols> {
let timing = crate::profile::enabled().then(Instant::now);
let ext = file.extension().and_then(|e| e.to_str())?;
let plugin = lang::plugin_for_extension(ext)?;
let rel = file
.strip_prefix(root)
.unwrap_or(file)
.to_string_lossy()
.into_owned();
let source = std::fs::read_to_string(file).ok()?;
if let Some(n) = needle
&& !contains_ascii_ci(source.as_bytes(), n)
{
return None;
}
let content_hash = content_hash(&source);
let symbols = plugin.extract(&rel, &source);
if let Some(started) = timing {
let elapsed = started.elapsed();
PARSE_US.fetch_add(
elapsed.as_micros() as u64,
std::sync::atomic::Ordering::Relaxed,
);
crate::profile::slow(elapsed, || rel.clone());
}
Some(crate::store::FileSymbols {
path: rel,
language: plugin.language().to_string(),
mtime: file_mtime(file),
content_hash,
generated: is_generated(&source),
symbols,
})
}
const HEADER_LINES: usize = 20;
pub(crate) fn is_generated(source: &str) -> bool {
source.lines().take(HEADER_LINES).any(|line| {
let line = line.trim_start();
let comment = line
.chars()
.next()
.is_some_and(|c| !c.is_alphanumeric() && !matches!(c, '"' | '\'' | '`'));
let lower = line.to_lowercase();
comment
&& (lower.contains("@generated")
|| (lower.contains("generated") && lower.contains("do not edit")))
})
}
fn past(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|d| Instant::now() >= d)
}
fn parse_files(
root: &Path,
paths: &[std::path::PathBuf],
deadline: Option<Instant>,
needle: Option<&[u8]>,
) -> (Vec<crate::store::FileSymbols>, bool) {
use std::sync::atomic::{AtomicBool, Ordering};
let workers = parse_jobs().min(paths.len());
if workers <= 1 {
let mut out = Vec::new();
for p in paths {
if past(deadline) {
return (out, false);
}
if let Some(parsed) = parse_file(root, p, needle) {
out.push(parsed);
}
}
return (out, true);
}
let bailed = AtomicBool::new(false);
let chunk_size = paths.len().div_ceil(workers);
let mut out = Vec::new();
std::thread::scope(|s| {
let handles: Vec<_> = paths
.chunks(chunk_size)
.map(|chunk| {
let bailed = &bailed;
s.spawn(move || {
let mut local = Vec::new();
for p in chunk {
if past(deadline) {
bailed.store(true, Ordering::Relaxed);
break;
}
if let Some(parsed) = parse_file(root, p, needle) {
local.push(parsed);
}
}
local
})
})
.collect();
for h in handles {
out.extend(h.join().unwrap_or_default());
}
});
(out, !bailed.load(Ordering::Relaxed))
}
fn capture_commit_times(store: &mut Store, repo_id: i64, root: &Path) {
let Some(head) = git_head(root) else { return };
let last = store.git_ts_head(repo_id).ok().flatten();
if last.as_deref() == Some(head.as_str()) {
return; }
let first = last.is_none();
let times = last
.and_then(|old| git_commit_times_range(root, &old, 1000))
.unwrap_or_else(|| git_commit_times(root, 1000));
if !times.is_empty() {
if store.set_file_git_ts(repo_id, ×).is_err() {
return; }
} else if first {
return; }
let _ = store.set_git_ts_head(repo_id, &head);
}
fn git_commit_times(root: &Path, limit: usize) -> HashMap<String, i64> {
match git_output(
root,
&[
"log",
&format!("-n{limit}"),
"--name-only",
"--pretty=format:%ct",
],
) {
Some(text) => parse_git_log(&text),
None => HashMap::new(),
}
}
fn git_commit_times_range(root: &Path, old: &str, limit: usize) -> Option<HashMap<String, i64>> {
git_output(
root,
&[
"log",
&format!("-n{limit}"),
"--name-only",
"--pretty=format:%ct",
&format!("{old}..HEAD"),
],
)
.map(|text| parse_git_log(&text))
}
fn parse_git_log(text: &str) -> HashMap<String, i64> {
let mut map = HashMap::new();
let mut current_ts = 0i64;
for line in text.lines() {
if line.is_empty() {
continue;
}
if let Ok(ts) = line.parse::<i64>() {
current_ts = ts;
} else {
map.entry(line.to_string()).or_insert(current_ts);
}
}
map
}
pub(crate) fn scan(
root: &Path,
skip: &HashSet<String>,
deadline: Option<Instant>,
needle: Option<&[u8]>,
) -> Vec<crate::store::FileSymbols> {
let needle = needle.filter(|n| !n.is_empty());
let candidates: Box<dyn Iterator<Item = std::path::PathBuf> + Send> =
match git_source_candidates(root).filter(|paths| !paths.is_empty()) {
Some(paths) => Box::new(paths.into_iter()),
None => Box::new(fs_walk_candidates(vec![root.to_path_buf()], deadline)),
};
let mut out: Vec<crate::store::FileSymbols> = Vec::new();
let keep = |rel: &str, _: &Path| !skip.contains(rel); let _ = stream_walk(
root,
candidates,
deadline,
None,
needle,
HashSet::new(),
keep,
None,
|fs| {
out.push(fs);
Ok(())
},
);
out
}
fn contains_ascii_ci(haystack: &[u8], needle: &[u8]) -> bool {
if needle.len() > haystack.len() {
return false;
}
haystack
.windows(needle.len())
.any(|w| w.eq_ignore_ascii_case(needle))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Refresh {
Unchanged,
Updated,
}
pub(crate) fn is_git_repo(root: &Path) -> bool {
repo_root(root).is_some()
}
pub(crate) fn repo_root(path: &Path) -> Option<std::path::PathBuf> {
let start = path.canonicalize().ok()?;
start
.ancestors()
.find(|a| a.join(".git").exists())
.map(Path::to_path_buf)
}
pub(crate) fn git_head(root: &Path) -> Option<String> {
let git_dir = root.join(".git");
if !git_dir.is_dir() {
return git_output(root, &["rev-parse", "HEAD"]);
}
let head = std::fs::read_to_string(git_dir.join("HEAD")).ok()?;
let head = head.trim();
let Some(git_ref) = head.strip_prefix("ref: ") else {
return (!head.is_empty()).then(|| head.to_string());
};
if let Ok(sha) = std::fs::read_to_string(git_dir.join(git_ref)) {
let sha = sha.trim();
if !sha.is_empty() {
return Some(sha.to_string());
}
}
let packed = std::fs::read_to_string(git_dir.join("packed-refs")).ok()?;
packed
.lines()
.find_map(|l| l.strip_suffix(&format!(" {git_ref}")))
.map(|sha| sha.trim().to_string())
.filter(|s| !s.is_empty())
}
const UNBORN_HEAD: &str = "unborn";
pub(crate) fn head_state(root: &Path) -> Option<String> {
git_head(root).or_else(|| is_git_repo(root).then(|| UNBORN_HEAD.to_string()))
}
pub(crate) fn dirty_files(root: &Path) -> Vec<String> {
Command::new("git")
.arg("-C")
.arg(root)
.args(["status", "--porcelain", "-z", "--untracked-files=no"])
.output()
.ok()
.filter(|o| o.status.success())
.map(|o| parse_porcelain_z(&o.stdout))
.unwrap_or_default()
}
fn parse_porcelain_z(out: &[u8]) -> Vec<String> {
let mut paths = Vec::new();
let mut entries = out.split(|&b| b == 0).filter(|e| e.len() > 3);
while let Some(entry) = entries.next() {
let (xy, path) = entry.split_at(3);
paths.push(String::from_utf8_lossy(path).into_owned());
if xy[..2].iter().any(|c| matches!(c, b'R' | b'C'))
&& let Some(source) = entries.next()
{
paths.push(String::from_utf8_lossy(source).into_owned());
}
}
paths
}
fn has_unindexed_edits(store: &Store, repository_id: i64, root: &Path, dirty: &[String]) -> bool {
dirty.iter().any(|rel| {
if !is_source(rel) {
return false;
}
let on_disk = file_mtime(&root.join(rel));
match store.file_mtime(repository_id, rel) {
Ok(Some(indexed)) => indexed.is_none() || indexed != on_disk,
Ok(None) => on_disk.is_some(),
Err(_) => true,
}
})
}
pub(crate) fn has_unindexed_changes(
store: &Store,
repository_id: i64,
root: &Path,
dirty: &[String],
) -> bool {
let prior = store.edited_files(repository_id).unwrap_or_default();
let unsettled: Vec<String> = prior
.iter()
.filter(|f| !dirty.contains(f))
.filter(|f| has_unindexed_edits(store, repository_id, root, std::slice::from_ref(f)))
.cloned()
.collect();
let changed = !unsettled.is_empty() || has_unindexed_edits(store, repository_id, root, dirty);
let mut edited: Vec<String> = dirty
.iter()
.filter(|f| is_source(f))
.cloned()
.chain(unsettled)
.collect();
edited.sort();
edited.dedup();
let mut prior = prior;
prior.sort();
if edited != prior {
let _ = store.set_edited_files(repository_id, &edited);
}
changed
}
fn is_source(rel: &str) -> bool {
let path = Path::new(rel);
let hidden = path.components().any(|c| match c {
std::path::Component::Normal(name) => name.as_encoded_bytes().starts_with(b"."),
_ => false,
});
!hidden
&& path
.extension()
.and_then(|e| e.to_str())
.is_some_and(|e| lang::plugin_for_extension(e).is_some())
}
pub(crate) fn branch_changed_files(root: &Path) -> Vec<String> {
let Some(branch) = head_branch(root) else {
return Vec::new();
};
if is_trunk(&branch) {
return Vec::new();
}
let Some(trunk) = trunk_ref(root) else {
return Vec::new();
};
let committed = {
let root = root.to_path_buf();
let spec = format!("{trunk}...HEAD");
std::thread::spawn(move || git_output(&root, &["diff", "--name-only", &spec]))
};
let working = git_output(root, &["diff", "--name-only", "HEAD"]);
let mut files: HashMap<String, ()> = HashMap::new();
for out in [committed.join().ok().flatten(), working]
.into_iter()
.flatten()
{
files.extend(
out.lines()
.filter(|l| !l.is_empty())
.map(|l| (l.to_string(), ())),
);
}
files.into_keys().collect()
}
pub(crate) fn branch_files_stamp(root: &Path) -> Option<String> {
let git_dir = root.join(".git");
if !git_dir.is_dir() {
return None;
}
let stamp = |name: &str| -> u64 {
std::fs::metadata(git_dir.join(name))
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_secs())
.unwrap_or(0)
};
Some(format!("{}:{}", stamp("HEAD"), stamp("index")))
}
pub(crate) fn git_state_stamp(root: &Path, head: &str) -> Option<String> {
let git_dir = root.join(".git");
if !git_dir.is_dir() || head_state(root)? != head {
return None;
}
let index = std::fs::metadata(git_dir.join("index"))
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map_or(0, |d| d.as_nanos());
Some(format!("{}\n{head}\n{index}", root.display()))
}
fn head_branch(root: &Path) -> Option<String> {
let git_dir = root.join(".git");
if !git_dir.is_dir() {
return is_git_repo(root)
.then(|| git_output(root, &["rev-parse", "--abbrev-ref", "HEAD"]))
.flatten();
}
let head = std::fs::read_to_string(git_dir.join("HEAD")).ok()?;
let branch = head.trim().strip_prefix("ref: refs/heads/")?;
(!branch.is_empty()).then(|| branch.to_string())
}
fn is_trunk(branch: &str) -> bool {
matches!(branch, "main" | "master" | "trunk")
}
fn trunk_ref(root: &Path) -> Option<String> {
let git_dir = root.join(".git");
if !git_dir.is_dir() {
return ["main", "master"]
.into_iter()
.find(|name| git_output(root, &["rev-parse", "--verify", "--quiet", name]).is_some())
.map(str::to_string);
}
let packed = std::fs::read_to_string(git_dir.join("packed-refs")).unwrap_or_default();
["main", "master"].into_iter().find_map(|name| {
let loose = git_dir.join("refs/heads").join(name).exists();
let is_packed = packed
.lines()
.any(|l| l.ends_with(&format!(" refs/heads/{name}")));
(loose || is_packed).then(|| name.to_string())
})
}
pub(crate) fn refresh_file(
store: &mut Store,
repository_id: i64,
root: &Path,
rel: &str,
) -> Result<Refresh, Box<dyn std::error::Error>> {
let path = root.join(rel);
let mtime = file_mtime(&path);
if mtime.is_some() && store.file_mtime(repository_id, rel)? == Some(mtime) {
return Ok(Refresh::Unchanged);
}
let source = match std::fs::read_to_string(&path) {
Ok(s) => s,
Err(_) => return Ok(Refresh::Unchanged), };
let hash = content_hash(&source);
if store.file_unchanged(repository_id, rel, &hash)? {
store.set_file_mtime(repository_id, rel, mtime)?;
return Ok(Refresh::Unchanged);
}
let ext = path
.extension()
.and_then(|e| e.to_str())
.unwrap_or_default();
let plugin = lang::plugin_for_extension(ext);
let symbols = match plugin {
Some(plugin) => plugin.extract(rel, &source),
None => Vec::new(),
};
let language = plugin.map_or("unknown", |p| p.language());
store.replace_files(
repository_id,
&[crate::store::FileSymbols {
path: rel.to_string(),
language: language.to_string(),
mtime,
content_hash: hash,
generated: is_generated(&source),
symbols,
}],
)?;
let _ = store.note_edited_file(repository_id, rel);
Ok(Refresh::Updated)
}
pub(crate) fn current_definitions(
store: &Store,
repository_id: Option<i64>,
identity: &str,
root: &Path,
rel: &str,
) -> Vec<crate::store::SymbolRow> {
let path = root.join(rel);
if let Some(repo_id) = repository_id
&& let Some(mtime) = file_mtime(&path)
&& store.file_mtime(repo_id, rel).ok() == Some(Some(Some(mtime)))
&& let Ok(rows) = store.symbols_in_file(repo_id, rel)
{
return rows;
}
let Some(plugin) = path
.extension()
.and_then(|e| e.to_str())
.and_then(lang::plugin_for_extension)
else {
return Vec::new();
};
let Ok(source) = std::fs::read_to_string(&path) else {
return Vec::new();
};
let generated = is_generated(&source);
plugin
.extract(rel, &source)
.into_iter()
.map(|s| crate::store::SymbolRow::live(s, repository_id.unwrap_or(-1), identity, generated))
.collect()
}
pub(crate) fn pushed_head(root: &Path) -> Option<String> {
let out = Command::new("git")
.arg("-C")
.arg(root)
.args(["rev-list", "--boundary", "HEAD", "--not", "--remotes", "--"])
.output()
.ok()?;
if !out.status.success() {
return None;
}
let text = String::from_utf8_lossy(&out.stdout);
if text.trim().is_empty() {
return git_head(root);
}
text.lines()
.find_map(|l| l.strip_prefix('-'))
.map(str::to_string)
}
pub(crate) fn detect_identity(root: &Path) -> RepoIdentity {
let remotes = if is_git_repo(root) {
&["origin", "upstream"][..]
} else {
&[]
};
for remote in remotes {
if let Some(url) = git_output(root, &["remote", "get-url", remote])
&& let Some(id) = RepoIdentity::from_remote_url(&url)
{
return id;
}
}
let abs = root.canonicalize().unwrap_or_else(|_| root.to_path_buf());
RepoIdentity::local(&abs.to_string_lossy())
}
fn git_output(root: &Path, args: &[&str]) -> Option<String> {
let out = Command::new("git")
.arg("-C")
.arg(root)
.args(args)
.output()
.ok()?;
if !out.status.success() {
return None;
}
let s = String::from_utf8(out.stdout).ok()?.trim().to_string();
if s.is_empty() { None } else { Some(s) }
}
fn content_hash(source: &str) -> String {
let mut hasher = std::collections::hash_map::DefaultHasher::new();
source.hash(&mut hasher);
format!("{:016x}", hasher.finish())
}
fn file_mtime(path: &Path) -> Option<i64> {
let modified = std::fs::metadata(path).ok()?.modified().ok()?;
let nanos = modified.duration_since(UNIX_EPOCH).ok()?.as_nanos();
Some(nanos as i64)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn generated_files_are_known_by_their_header() {
for src in [
"// Code generated by \"stringer -type=Kind\"; DO NOT EDIT.\n\npackage kinds\n",
"// Copyright 2024\n// License: MIT\n\n// Code generated by protoc-gen-go. DO NOT EDIT.\npackage pb\n",
"# Generated by the protocol buffer compiler. DO NOT EDIT!\nimport sys\n",
"/**\n * @generated\n */\nexport type Foo = {};\n",
] {
assert!(is_generated(src), "generated: {src:?}");
}
for src in [
"package kinds\n\nfunc String() {}\n",
"const HEADER = \"// Code generated by gen; DO NOT EDIT.\";\n",
"fmt.Println(\"// Code generated; DO NOT EDIT.\")\n",
"// Regenerate the docs with `make docs`.\n",
] {
assert!(!is_generated(src), "hand-written: {src:?}");
}
let late = format!(
"{}// Code generated; DO NOT EDIT.\n",
"x := 1\n".repeat(HEADER_LINES)
);
assert!(!is_generated(&late), "only the header counts");
}
#[test]
fn a_filesystem_walk_stops_descending_at_its_deadline() {
let dir = std::env::temp_dir().join(format!("rq-walk-deadline-{}", std::process::id()));
std::fs::create_dir_all(dir.join("a/b")).unwrap();
std::fs::write(dir.join("a/b/x.rb"), "class X\nend\n").unwrap();
let files = |deadline| fs_walk_candidates(vec![dir.clone()], deadline).count();
assert_eq!(files(None), 1);
assert_eq!(files(Some(Instant::now())), 0);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn sweep_outcome_guards_against_a_failed_empty_walk() {
assert_eq!(
sweep_outcome(true, true, false, true, true),
(true, "complete")
);
assert_eq!(
sweep_outcome(true, true, true, false, true),
(true, "complete")
);
assert_eq!(
sweep_outcome(true, true, true, true, true),
(false, "warming")
);
assert_eq!(
sweep_outcome(true, true, true, true, false),
(true, "complete")
);
assert_eq!(
sweep_outcome(false, true, false, true, true),
(false, "warming")
);
assert_eq!(
sweep_outcome(true, false, false, true, true),
(false, "warming")
);
}
#[test]
fn parse_jobs_defaults_to_one_per_available_core() {
set_parse_jobs(0); let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1);
if std::env::var_os("RQ_JOBS").is_none() {
assert_eq!(parse_jobs(), cores);
}
set_parse_jobs(3);
assert_eq!(parse_jobs(), 3, "an explicit --jobs still wins");
set_parse_jobs(0);
}
#[test]
fn content_hash_is_stable_and_distinguishes() {
assert_eq!(
content_hash("class Foo\nend"),
content_hash("class Foo\nend")
);
assert_ne!(
content_hash("class Foo\nend"),
content_hash("class Bar\nend")
);
}
#[test]
fn porcelain_z_yields_every_path_including_a_rename_source() {
let out = b" M a.rb\0M lib/b.rb\0R new.rb\0old.rb\0D gone.rb\0";
assert_eq!(
parse_porcelain_z(out),
["a.rb", "lib/b.rb", "new.rb", "old.rb", "gone.rb"]
);
assert!(parse_porcelain_z(b"").is_empty());
}
#[test]
fn trunk_names_are_recognized() {
assert!(is_trunk("main"));
assert!(is_trunk("master"));
assert!(!is_trunk("feature/x"));
assert!(!is_trunk("dpep/fix"));
}
#[test]
fn prioritize_by_path_is_loose_but_targeted() {
let root = Path::new("/repo");
let paths: Vec<std::path::PathBuf> = [
"companies.rb", "app/employee.rb", "lib/EmpController.rb", "employers.rb", "app/controllers/x.rb", ]
.iter()
.map(|p| root.join(p))
.collect();
let out = prioritize_by_path(paths.clone(), root, Some("employeescontroller"));
let name = |p: &std::path::PathBuf| p.file_name().unwrap().to_str().unwrap().to_string();
let front: Vec<String> = out[..3].iter().map(name).collect();
assert!(front.contains(&"employee.rb".to_string()), "{front:?}");
assert!(front.contains(&"EmpController.rb".to_string()), "{front:?}");
assert!(front.contains(&"employers.rb".to_string()), "{front:?}");
let tail: Vec<String> = out[3..].iter().map(name).collect();
assert!(tail.contains(&"companies.rb".to_string()), "{tail:?}");
assert!(tail.contains(&"x.rb".to_string()), "{tail:?}"); assert_eq!(prioritize_by_path(paths.clone(), root, None), paths);
}
#[test]
fn detects_git_work_tree_natively() {
let dir = std::env::temp_dir().join(format!("rq-reporoot-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(dir.join("sub")).unwrap();
assert!(!is_git_repo(&dir), "no .git yet");
std::fs::create_dir_all(dir.join(".git")).unwrap();
assert!(is_git_repo(&dir), "a .git entry marks a work tree");
assert_eq!(
repo_root(&dir.join("sub")).unwrap(),
dir.canonicalize().unwrap()
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn parses_git_log_keeping_most_recent_commit_per_file() {
let log = "1700000000\n\na.rb\nb.rb\n1699990000\n\na.rb\nc.rb\n";
let map = parse_git_log(log);
assert_eq!(map.get("a.rb"), Some(&1700000000));
assert_eq!(map.get("b.rb"), Some(&1700000000));
assert_eq!(map.get("c.rb"), Some(&1699990000));
assert_eq!(map.len(), 3);
}
}