use crate::structure;
use crate::CliError;
use core_api::{
default_max_edges, Direction, GraphError, IngestOptions, Predicate, ResultSet, RuleDef,
SharedDb, Value, WriteGuard, WRITE_LOCK_WAIT,
};
use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::process::Command;
pub const DEFAULT_MAX_COMMITS_PER_FILE: usize = 200;
pub use core_api::repograph::DEFAULT_EXCLUDES;
const CO_CHANGE_MIN: f64 = 0.25;
pub const SNAPSHOT_WAL_BYTES: u64 = 4 * 1024 * 1024;
pub(crate) const SYNC_KEY: &str = "__mushroomdb_git_sync__";
pub const BUSY_MESSAGE: &str = "another mushroomdb process is writing; retry";
const PR_FK_RULE: &str = "auto_fk_commit_pr_id";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngestGitOpts {
pub repo: PathBuf,
pub exclude: Vec<String>,
pub max_commits_per_file: usize,
pub recurse_submodules: bool,
pub prs: bool,
pub structure: bool,
pub docs: bool,
pub ensure_gitignore: bool,
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct IngestGitReport {
pub commits: usize,
pub files: usize,
pub authors: usize,
pub renamed: usize,
pub deleted: usize,
pub incremental: bool,
pub rules_created: Vec<String>,
pub submodules: usize,
pub prs: usize,
pub gitignore_added: bool,
pub structure: crate::structure::StructureReport,
}
#[derive(Default)]
struct StructureWork {
touched: BTreeSet<String>,
stale: BTreeSet<String>,
}
#[derive(Debug, Clone)]
struct RepoUnit {
path: PathBuf,
prefix: String,
sync_key: String,
}
#[derive(Debug)]
enum Change {
Added(String),
Modified(String),
Deleted(String),
Renamed { from: String, to: String },
}
#[derive(Debug)]
struct GitCommit {
sha: String,
author_name: String,
author_email: String,
ts: i64,
subject: String,
changes: Vec<Change>,
}
fn excluded(path: &str, patterns: &[String]) -> bool {
core_api::repograph::path_excluded(path, patterns)
}
fn git_output(repo: &Path, args: &[&str]) -> Result<std::process::Output, CliError> {
Command::new("git")
.arg("-C")
.arg(repo)
.args(args)
.output()
.map_err(|e| CliError(format!("cannot run git in {}: {e}", repo.display())))
}
fn read_log(repo: &Path, since: Option<&str>, head: &str) -> Result<Vec<GitCommit>, CliError> {
if let Some(s) = since {
let spec = format!("{s}^{{commit}}");
if !git_output(repo, &["cat-file", "-e", &spec])?
.status
.success()
{
return Err(CliError(format!(
"recorded sync head {s} is not in {} (history rewritten?); \
ingest into a fresh database directory",
repo.display()
)));
}
}
let mut cmd = Command::new("git");
cmd.arg("-C").arg(repo).args([
"-c",
"core.quotePath=false",
"log",
"--reverse",
"--name-status",
"-M",
"--no-color",
"--format=%x1e%H%x1f%aN%x1f%aE%x1f%at%x1f%s",
]);
match since {
Some(s) => cmd.arg(format!("{s}..{head}")),
None => cmd.arg(head),
};
let out = cmd
.output()
.map_err(|e| CliError(format!("cannot run git: {e}")))?;
if !out.status.success() {
return Err(CliError(format!(
"git log failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
)));
}
let text = String::from_utf8_lossy(&out.stdout);
let mut commits = Vec::new();
for block in text.split('\x1e').filter(|b| !b.trim().is_empty()) {
let mut lines = block.lines();
let header = lines.next().unwrap_or("");
let f: Vec<&str> = header.split('\x1f').collect();
if f.len() < 5 {
continue;
}
let mut changes = Vec::new();
for l in lines {
let cols: Vec<&str> = l.split('\t').collect();
match cols.as_slice() {
[s, p] if s.starts_with('A') => changes.push(Change::Added((*p).to_string())),
[s, p] if s.starts_with('M') || s.starts_with('T') => {
changes.push(Change::Modified((*p).to_string()))
}
[s, p] if s.starts_with('D') => changes.push(Change::Deleted((*p).to_string())),
[s, from, to] if s.starts_with('R') => changes.push(Change::Renamed {
from: (*from).to_string(),
to: (*to).to_string(),
}),
[s, _from, to] if s.starts_with('C') => {
changes.push(Change::Added((*to).to_string()))
}
_ => {}
}
}
commits.push(GitCommit {
sha: f[0].into(),
author_name: f[1].into(),
author_email: f[2].into(),
ts: f[3].parse().unwrap_or(0),
subject: f[4].into(),
changes,
});
}
Ok(commits)
}
fn head_sha(repo: &Path) -> Result<Option<String>, CliError> {
let out = git_output(repo, &["rev-parse", "--verify", "-q", "HEAD^{commit}"])?;
if !out.status.success() {
if !git_output(repo, &["rev-parse", "--git-dir"])?
.status
.success()
{
return Err(CliError(format!(
"not a git repository: {}",
repo.display()
)));
}
return Ok(None);
}
let sha = String::from_utf8_lossy(&out.stdout).trim().to_string();
if sha.is_empty() {
return Err(CliError(format!(
"git could not resolve HEAD in {}",
repo.display()
)));
}
Ok(Some(sha))
}
fn canonical(p: &Path) -> PathBuf {
std::fs::canonicalize(p).unwrap_or_else(|_| p.to_path_buf())
}
fn gitlink_paths(unit: &Path) -> BTreeSet<String> {
let out = Command::new("git")
.arg("config")
.arg("--file")
.arg(unit.join(".gitmodules"))
.args(["--get-regexp", r"^submodule\..*\.path$"])
.output();
let Ok(out) = out else { return BTreeSet::new() };
if !out.status.success() {
return BTreeSet::new(); }
String::from_utf8_lossy(&out.stdout)
.lines()
.filter_map(|l| l.split_once(' '))
.map(|(_, path)| path.trim().to_string())
.filter(|p| !p.is_empty())
.collect()
}
fn repo_units(repo: &Path, recurse: bool) -> Result<Vec<RepoUnit>, CliError> {
let mut units = vec![RepoUnit {
path: repo.to_path_buf(),
prefix: String::new(),
sync_key: SYNC_KEY.to_string(),
}];
if !recurse {
return Ok(units);
}
let out = git_output(
repo,
&[
"submodule",
"foreach",
"--quiet",
"--recursive",
"echo \"$displaypath\"",
],
)?;
if !out.status.success() {
return Ok(units);
}
let mut paths: Vec<String> = String::from_utf8_lossy(&out.stdout)
.lines()
.map(|l| l.trim_end_matches('/').trim().to_string())
.filter(|l| !l.is_empty())
.collect();
paths.sort();
paths.dedup();
for dp in paths {
let path = repo.join(&dp);
if !git_output(&path, &["rev-parse", "--git-dir"])?
.status
.success()
{
continue;
}
units.push(RepoUnit {
prefix: format!("{dp}/"),
sync_key: format!("{SYNC_KEY}:{dp}"),
path,
});
}
Ok(units)
}
fn localise(log: &mut [GitCommit], prefix: &str, gitlinks: &BTreeSet<String>) {
let keep = |p: &String| !gitlinks.contains(p.as_str());
for c in log.iter_mut() {
c.changes.retain(|ch| match ch {
Change::Added(p) | Change::Modified(p) | Change::Deleted(p) => keep(p),
Change::Renamed { from, to } => keep(from) && keep(to),
});
if prefix.is_empty() {
continue;
}
for ch in c.changes.iter_mut() {
match ch {
Change::Added(p) | Change::Modified(p) | Change::Deleted(p) => {
*p = format!("{prefix}{p}")
}
Change::Renamed { from, to } => {
*from = format!("{prefix}{from}");
*to = format!("{prefix}{to}");
}
}
}
}
}
fn ensure_gitignore(repo: &Path, db_dir: &Path) -> Result<bool, CliError> {
let Ok(rel) = canonical(db_dir)
.strip_prefix(repo)
.map(|p| p.to_path_buf())
else {
return Ok(false);
};
if rel.as_os_str().is_empty() {
return Ok(false);
}
let line = format!("{}/", rel.to_string_lossy().replace('\\', "/"));
let path = repo.join(".gitignore");
let current = match std::fs::read_to_string(&path) {
Ok(s) => s,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => String::new(),
Err(e) => return Err(CliError(format!("cannot read {}: {e}", path.display()))),
};
let bare = line.trim_end_matches('/');
if current
.lines()
.map(|l| l.trim())
.any(|l| l == line || l == bare || l == format!("/{line}") || l == format!("/{bare}"))
{
return Ok(false);
}
let mut next = current;
if !next.is_empty() && !next.ends_with('\n') {
next.push('\n');
}
next.push_str(&line);
next.push('\n');
std::fs::write(&path, next)
.map_err(|e| CliError(format!("cannot write {}: {e}", path.display())))?;
Ok(true)
}
const AUTHOR_COUNT_SEP: char = '\t';
#[derive(Default, Clone, Debug, PartialEq, Eq)]
struct AuthorState {
names: Vec<(String, usize)>,
}
impl AuthorState {
fn touch(&mut self, name: &str) {
self.add(name, 1);
}
fn add(&mut self, name: &str, n: usize) {
match self.names.iter_mut().find(|(held, _)| held == name) {
Some(entry) => entry.1 += n,
None => self.names.push((name.to_string(), n)),
}
}
fn absorb(&mut self, other: &AuthorState) {
for (name, n) in &other.names {
self.add(name, *n);
}
}
fn display_name(&self) -> String {
let mut best: Option<&(String, usize)> = None;
for entry in &self.names {
if best.is_none_or(|(_, n)| entry.1 > *n) {
best = Some(entry);
}
}
best.map(|(name, _)| name.clone()).unwrap_or_default()
}
fn name_counts_value(&self) -> Value {
Value::List(
self.names
.iter()
.map(|(name, n)| Value::Str(format!("{name}{AUTHOR_COUNT_SEP}{n}")))
.collect(),
)
}
fn set_name_counts(&mut self, list: &[Value]) {
for v in list {
let Value::Str(s) = v else { continue };
let Some((name, n)) = s.rsplit_once(AUTHOR_COUNT_SEP) else {
continue;
};
let Ok(n) = n.parse::<usize>() else { continue };
if !name.is_empty() {
self.add(name, n);
}
}
}
}
fn file_props(st: &FileState, path: &str) -> Vec<(String, Value)> {
let commits = &st.commits;
let dir = path
.rsplit_once('/')
.map(|(d, _)| d)
.unwrap_or("")
.to_string();
let ext = path
.rsplit_once('.')
.map(|(_, e)| e)
.unwrap_or("")
.to_string();
vec![
("id".into(), Value::Str(path.into())),
("path".into(), Value::Str(path.into())),
("dir".into(), Value::Str(dir)),
("ext".into(), Value::Str(ext)),
(
"commits".into(),
Value::List(commits.iter().map(|s| Value::Str(s.clone())).collect()),
),
("n_commits".into(), Value::Int(st.n_commits as i64)),
("top_author_id".into(), Value::Str(st.top_author())),
("author_counts".into(), st.author_counts_value()),
]
}
#[derive(Default, Clone)]
struct FileState {
commits: Vec<String>,
by_author: BTreeMap<String, usize>,
n_commits: usize,
}
impl FileState {
fn touch(&mut self, sha: &str, author: &str, cap: usize) {
self.commits.push(sha.to_string());
if cap > 0 && self.commits.len() > cap {
self.commits.remove(0);
}
self.n_commits += 1;
*self.by_author.entry(author.to_string()).or_default() += 1;
}
fn top_author(&self) -> String {
self.by_author
.iter()
.max_by(|a, b| a.1.cmp(b.1).then(b.0.cmp(a.0)))
.map(|(a, _)| a.clone())
.unwrap_or_default()
}
fn author_counts_value(&self) -> Value {
Value::List(
self.by_author
.iter()
.map(|(email, n)| Value::Str(format!("{email}{AUTHOR_COUNT_SEP}{n}")))
.collect(),
)
}
fn set_author_counts(&mut self, list: &[Value]) {
for v in list {
let Value::Str(s) = v else { continue };
let Some((email, n)) = s.rsplit_once(AUTHOR_COUNT_SEP) else {
continue;
};
let Ok(n) = n.parse::<usize>() else { continue };
if !email.is_empty() {
*self.by_author.entry(email.to_string()).or_default() += n;
}
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PullRequest {
number: i64,
title: String,
url: String,
merged_at: String,
author_login: String,
merge_sha: Option<String>,
}
fn pr_key(number: i64) -> String {
format!("pr:{number}")
}
fn subject_pr(subject: &str) -> Option<i64> {
let rest = subject.strip_suffix(')')?;
let at = rest.rfind("(#")?;
let digits = &rest[at + 2..];
if digits.is_empty() || !digits.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
digits.parse().ok()
}
fn fetch_prs(repo: &Path) -> Vec<PullRequest> {
let out = Command::new("gh")
.current_dir(repo)
.args([
"pr",
"list",
"--state",
"merged",
"--limit",
"1000",
"--json",
"number,title,url,mergedAt,mergeCommit,author",
])
.output();
let out = match out {
Ok(o) => o,
Err(_) => {
eprintln!("ingest-git: --prs skipped: gh is not on PATH");
return Vec::new();
}
};
if !out.status.success() {
let detail = String::from_utf8_lossy(&out.stderr)
.lines()
.find(|l| !l.trim().is_empty())
.unwrap_or("no detail")
.to_string();
eprintln!("ingest-git: --prs skipped: gh pr list failed: {detail}");
return Vec::new();
}
let parsed: serde_json::Value = match serde_json::from_slice(&out.stdout) {
Ok(v) => v,
Err(e) => {
eprintln!("ingest-git: --prs skipped: gh pr list output is not JSON: {e}");
return Vec::new();
}
};
let Some(items) = parsed.as_array() else {
eprintln!("ingest-git: --prs skipped: gh pr list did not return a list");
return Vec::new();
};
let mut prs: Vec<PullRequest> = items
.iter()
.filter_map(|v| {
let number = v.get("number")?.as_i64()?;
let str_at = |k: &str| {
v.get(k)
.and_then(|x| x.as_str())
.unwrap_or_default()
.to_string()
};
Some(PullRequest {
number,
title: str_at("title"),
url: str_at("url"),
merged_at: str_at("mergedAt"),
author_login: v
.get("author")
.and_then(|a| a.get("login"))
.and_then(|l| l.as_str())
.unwrap_or_default()
.to_string(),
merge_sha: v
.get("mergeCommit")
.and_then(|m| m.get("oid"))
.and_then(|o| o.as_str())
.filter(|s| !s.is_empty())
.map(str::to_string),
})
})
.collect();
prs.sort_by_key(|p| p.number);
prs.dedup_by_key(|p| p.number);
prs
}
fn ingest_prs(
w: &mut WriteGuard<'_>,
prs: &[PullRequest],
ingest: &IngestOptions,
report: &mut IngestGitReport,
) -> Result<(), CliError> {
let rows: Vec<BTreeMap<String, Value>> = prs
.iter()
.filter(|p| !w.has_node(&pr_key(p.number)))
.map(|p| {
BTreeMap::from([
("id".to_string(), Value::Str(pr_key(p.number))),
("number".to_string(), Value::Int(p.number)),
("title".to_string(), Value::Str(p.title.clone())),
("url".to_string(), Value::Str(p.url.clone())),
("merged_at".to_string(), Value::Str(p.merged_at.clone())),
(
"author_login".to_string(),
Value::Str(p.author_login.clone()),
),
])
})
.collect();
if !rows.is_empty() {
let r = w.ingest_with_edges("PR", rows, ingest, &[])?;
report.prs = r.inserted;
report.rules_created.extend(r.rules_created);
}
if !w.rules().iter().any(|r| r.name == PR_FK_RULE) {
let predicate = Predicate::KeyMatch {
field: "pr_id".into(),
};
let max_edges = Some(default_max_edges(&predicate));
w.create_rule(RuleDef {
name: PR_FK_RULE.into(),
src_label: "Commit".into(),
dst_label: "PR".into(),
predicate,
edge_type: "PR".into(),
weight_prop: None,
max_edges,
approximate: false,
via_label: None,
via_edge: None,
via_dir: None,
namespace: None,
})?;
report.rules_created.push(PR_FK_RULE.to_string());
}
if !w
.fulltext_pairs()
.contains(&("PR".to_string(), "title".to_string()))
{
w.enable_fulltext("PR", "title")?;
}
Ok(())
}
fn link_prs(
w: &mut WriteGuard<'_>,
prs: &[PullRequest],
ingest: &IngestOptions,
) -> Result<(), CliError> {
let by_sha: BTreeMap<&str, i64> = prs
.iter()
.filter_map(|p| p.merge_sha.as_deref().map(|s| (s, p.number)))
.collect();
let known: BTreeSet<i64> = prs.iter().map(|p| p.number).collect();
let rs = w.query(
"MATCH (c:Commit) RETURN c.id AS id, c.message AS message, c.pr_id AS pr_id",
&BTreeMap::new(),
)?;
let mut updates: Vec<(String, String)> = Vec::new();
let mut links: BTreeMap<i64, BTreeSet<String>> = BTreeMap::new();
for i in 0..rs.len() {
let Some(Value::Str(sha)) = rs.get(i, "id") else {
continue;
};
let subject = match rs.get(i, "message") {
Some(Value::Str(s)) => s.as_str(),
_ => "",
};
let Some(number) = by_sha
.get(sha.as_str())
.copied()
.or_else(|| subject_pr(subject).filter(|n| known.contains(n)))
else {
continue;
};
let key = pr_key(number);
links.entry(number).or_default().insert(sha.clone());
if rs.get(i, "pr_id") != Some(&Value::Str(key.clone())) {
updates.push((sha.clone(), key));
}
}
updates.sort(); for (sha, key) in updates {
w.set_prop(&sha, "pr_id", Value::Str(key))?;
}
let mut edges: Vec<(String, String, String)> = Vec::new();
for (number, shas) in links {
let src = pr_key(number);
let existing: BTreeSet<String> = w
.neighbors(&src, "MERGED_AS", Direction::Out)
.unwrap_or_default()
.into_iter()
.collect();
for sha in shas {
if !existing.contains(&sha) {
edges.push(("MERGED_AS".to_string(), src.clone(), sha));
}
}
}
if !edges.is_empty() {
w.ingest_with_edges("PR", Vec::new(), ingest, &edges)?;
}
Ok(())
}
const FILE_STATE_QUERY: &str =
"MATCH (f:File) WHERE startsWith(f.id, $prefix) RETURN f.id AS id, f.commits AS commits, \
f.n_commits AS n, f.top_author_id AS top, f.author_counts AS author_counts";
fn file_state_from(rs: &ResultSet, nested: &[String]) -> BTreeMap<String, FileState> {
let mut files = BTreeMap::new();
for i in 0..rs.len() {
let id = match rs.get(i, "id") {
Some(Value::Str(s)) => s.clone(),
_ => continue,
};
if nested.iter().any(|p| id.starts_with(p.as_str())) {
continue;
}
let mut st = FileState::default();
if let Some(Value::List(l)) = rs.get(i, "commits") {
st.commits = l
.iter()
.filter_map(|v| match v {
Value::Str(s) => Some(s.clone()),
_ => None,
})
.collect();
}
if let Some(Value::Int(n)) = rs.get(i, "n") {
st.n_commits = *n as usize;
}
match rs.get(i, "author_counts") {
Some(Value::List(l)) => st.set_author_counts(l),
_ => {
if let Some(Value::Str(t)) = rs.get(i, "top") {
st.by_author.insert(t.clone(), st.n_commits.max(1));
}
}
}
files.insert(id, st);
}
files
}
fn legacy_name_counts(
w: &WriteGuard<'_>,
p: &Pending,
) -> Result<BTreeMap<String, AuthorState>, CliError> {
let mut out: BTreeMap<String, AuthorState> = BTreeMap::new();
let stale = p.log.iter().any(|c| {
w.has_node(&c.author_email)
&& w.node_ref(&c.author_email)
.and_then(|n| n.prop("name_counts"))
.is_none()
});
let Some(head) = p.head.as_deref().filter(|_| stale) else {
return Ok(out);
};
for c in read_log(&p.unit.path, None, head)? {
out.entry(c.author_email).or_default().touch(&c.author_name);
}
Ok(out)
}
#[derive(Default)]
struct Walk {
files: BTreeMap<String, FileState>,
authors: BTreeMap<String, AuthorState>,
commit_rows: Vec<BTreeMap<String, Value>>,
touched_edges: Vec<(String, String, String)>,
dirty: BTreeSet<String>,
deleted: BTreeSet<String>,
renamed: Vec<(String, String)>,
alias: BTreeMap<String, String>,
}
impl Walk {
fn rename(&mut self, from: &str, to: &str, node_exists: bool) {
if let Some(e) = self.renamed.iter_mut().find(|(_, t)| t == from) {
e.1 = to.to_string();
} else if node_exists {
self.deleted.remove(from);
self.renamed.push((from.to_string(), to.to_string()));
} else {
self.deleted.insert(from.to_string());
}
for v in self.alias.values_mut() {
if v == from {
*v = to.to_string();
}
}
self.alias.insert(from.to_string(), to.to_string());
}
}
fn marker_flag_props(unit: &RepoUnit, opts: &IngestGitOpts) -> Vec<(String, Value)> {
vec![
(
"repo".into(),
Value::Str(unit.path.to_string_lossy().into_owned()),
),
("recurse".into(), Value::Bool(opts.recurse_submodules)),
("prs".into(), Value::Bool(opts.prs)),
("structure".into(), Value::Bool(opts.structure)),
("docs".into(), Value::Bool(opts.docs)),
]
}
struct Pending {
unit: RepoUnit,
since: Option<String>,
stale: bool,
head: Option<String>,
log: Vec<GitCommit>,
nested: Vec<String>,
}
fn snapshot_due(db_dir: &Path, full: bool) -> bool {
if full || !db_dir.join("snapshot.bin").exists() {
return true;
}
std::fs::metadata(db_dir.join("wal.bin")).is_ok_and(|m| m.len() > SNAPSHOT_WAL_BYTES)
}
fn snapshot_if_due(db: &SharedDb, db_dir: &Path, full: bool) {
if !snapshot_due(db_dir, full) {
return;
}
let taken = db
.write_with_wait(WRITE_LOCK_WAIT)
.and_then(|mut g| crate::snapshot_automatically(&mut g));
match taken {
Ok(()) | Err(GraphError::Busy { .. }) => {}
Err(e) => eprintln!("snapshot skipped: {e}"),
}
}
pub fn run_ingest_git(db_dir: &Path, opts: &IngestGitOpts) -> Result<IngestGitReport, CliError> {
let repo = canonical(&opts.repo);
let units = repo_units(&repo, opts.recurse_submodules)?;
let prefixes: Vec<String> = units
.iter()
.map(|u| u.prefix.clone())
.filter(|p| !p.is_empty())
.collect();
let db = SharedDb::open(db_dir)?;
let mut report = IngestGitReport {
submodules: units.len() - 1,
..Default::default()
};
if opts.ensure_gitignore {
report.gitignore_added = ensure_gitignore(&repo, db_dir)?;
}
let mut pending: Vec<Pending> = Vec::new();
{
let r = db.read();
for unit in units {
let nested = prefixes
.iter()
.filter(|p| **p != unit.prefix && p.starts_with(&unit.prefix))
.cloned()
.collect();
let (since, stale) = match r.node_ref(&unit.sync_key) {
Some(n) => (
match n.prop("sha") {
Some(Value::Str(s)) if !s.is_empty() => Some(s),
_ => None,
},
marker_flag_props(&unit, opts)
.into_iter()
.any(|(k, v)| n.prop(&k) != Some(v)),
),
None => (None, false),
};
pending.push(Pending {
unit,
since,
stale,
head: None,
log: Vec::new(),
nested,
});
}
}
report.incremental = pending[0].since.is_some();
for p in &mut pending {
let Some(head) = head_sha(&p.unit.path)? else {
continue; };
let mut log = read_log(&p.unit.path, p.since.as_deref(), &head)?;
localise(&mut log, &p.unit.prefix, &gitlink_paths(&p.unit.path));
p.head = Some(head);
p.log = log;
}
let prs = if opts.prs {
fetch_prs(&repo)
} else {
Vec::new()
};
if prs.is_empty() && pending.iter().all(|p| p.log.is_empty() && !p.stale) {
return Ok(report);
}
let mut w = db.write_with_wait(WRITE_LOCK_WAIT).map_err(|e| match e {
GraphError::Busy { .. } => CliError(BUSY_MESSAGE.to_string()),
other => CliError(other.to_string()),
})?;
let ingest = IngestOptions::default();
if !prs.is_empty() {
ingest_prs(&mut w, &prs, &ingest, &mut report)?;
}
let mut authors: BTreeSet<String> = BTreeSet::new();
let mut work = StructureWork::default();
for p in &pending {
if !p.log.is_empty() {
ingest_unit(
&mut w,
p,
opts,
&ingest,
&mut report,
&mut authors,
&mut work,
)?;
}
}
report.authors = authors.len();
if report.commits > 0 && !w.rules().iter().any(|r| r.name == "co_changed") {
let co = Predicate::Overlap {
field: "commits".into(),
min: CO_CHANGE_MIN,
};
w.create_rule(RuleDef {
name: "co_changed".into(),
src_label: "File".into(),
dst_label: "File".into(),
predicate: co.clone(),
edge_type: "CO_CHANGED".into(),
weight_prop: Some("score".into()),
max_edges: Some(10),
approximate: false,
via_label: None,
via_edge: None,
via_dir: None,
namespace: None,
})?;
w.create_rule(RuleDef {
name: "knows".into(),
src_label: "Author".into(),
dst_label: "File".into(),
predicate: co,
edge_type: "KNOWS".into(),
weight_prop: Some("score".into()),
max_edges: Some(20),
approximate: false,
via_label: Some("File".into()),
via_edge: Some("TOP_AUTHOR".into()),
via_dir: Some(Direction::In),
namespace: None,
})?;
report
.rules_created
.extend(["co_changed".to_string(), "knows".to_string()]);
for (l, field) in [("File", "path"), ("Commit", "message"), ("Author", "name")] {
if !w
.fulltext_pairs()
.contains(&(l.to_string(), field.to_string()))
{
w.enable_fulltext(l, field)?;
}
}
}
if opts.structure {
let full = !report.incremental || pending.iter().any(|p| p.stale);
report.structure = if full {
structure::refresh_all(&mut w, &repo, "", opts.docs)?
} else {
let mut paths = work.touched;
paths.extend(structure::importers_of(&w, &work.stale)?);
paths.extend(work.stale.iter().cloned());
let paths: Vec<String> = paths.into_iter().collect();
structure::refresh_files(&mut w, &repo, "", &paths, opts.docs)?
};
report
.rules_created
.extend(structure::ensure_rules_and_fulltext(&mut w)?);
}
if !prs.is_empty() {
link_prs(&mut w, &prs, &ingest)?;
}
for p in &pending {
write_marker(&mut w, p, opts)?;
}
drop(w);
snapshot_if_due(&db, db_dir, !report.incremental);
Ok(report)
}
fn now_unix() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
}
pub const SYNCED_AT: &str = "synced_at";
fn write_marker(w: &mut WriteGuard<'_>, p: &Pending, opts: &IngestGitOpts) -> Result<(), CliError> {
let Some(head) = p.head.as_deref() else {
return Ok(()); };
let mut props = marker_flag_props(&p.unit, opts);
if !p.log.is_empty() {
props.push(("sha".into(), Value::Str(head.to_string())));
}
let key = p.unit.sync_key.clone();
if w.has_node(&key) {
let mut changed = false;
for (k, v) in props {
let current = w.node_ref(&key).and_then(|n| n.prop(&k));
if current.as_ref() != Some(&v) {
w.set_prop(&key, &k, v)?;
changed = true;
}
}
if changed {
w.set_prop(&key, SYNCED_AT, Value::Int(now_unix()))?;
}
return Ok(());
}
if p.log.is_empty() {
return Ok(());
}
props.push(("id".into(), Value::Str(key.clone())));
props.push((SYNCED_AT.into(), Value::Int(now_unix())));
props.sort_by(|a, b| a.0.cmp(&b.0));
w.insert_node("GitSync", &key, props)?;
Ok(())
}
fn ingest_unit(
w: &mut WriteGuard<'_>,
p: &Pending,
opts: &IngestGitOpts,
ingest: &IngestOptions,
report: &mut IngestGitReport,
authors: &mut BTreeSet<String>,
work: &mut StructureWork,
) -> Result<(), CliError> {
let log = &p.log;
let incremental = p.since.is_some();
let mut walk = Walk {
files: if incremental {
let params =
BTreeMap::from([("prefix".to_string(), Value::Str(p.unit.prefix.clone()))]);
file_state_from(&w.query(FILE_STATE_QUERY, ¶ms)?, &p.nested)
} else {
BTreeMap::new()
},
..Default::default()
};
for c in log {
walk.authors
.entry(c.author_email.clone())
.or_default()
.touch(&c.author_name);
walk.commit_rows.push(BTreeMap::from([
("id".to_string(), Value::Str(c.sha.clone())),
("message".to_string(), Value::Str(c.subject.clone())),
("ts".to_string(), Value::Int(c.ts)),
("author_id".to_string(), Value::Str(c.author_email.clone())),
]));
for ch in &c.changes {
match ch {
Change::Added(p) | Change::Modified(p) => {
if excluded(p, &opts.exclude) {
continue;
}
walk.deleted.remove(p);
walk.files.entry(p.clone()).or_default().touch(
&c.sha,
&c.author_email,
opts.max_commits_per_file,
);
walk.dirty.insert(p.clone());
walk.touched_edges
.push(("TOUCHED".into(), c.sha.clone(), p.clone()));
}
Change::Deleted(p) => {
if excluded(p, &opts.exclude) {
continue;
}
walk.files.remove(p);
walk.dirty.remove(p);
walk.deleted.insert(p.clone());
}
Change::Renamed { from, to } => {
if excluded(to, &opts.exclude) {
walk.files.remove(from);
walk.dirty.remove(from);
walk.deleted.insert(from.clone());
continue;
}
let mut st = walk.files.remove(from).unwrap_or_default();
st.touch(&c.sha, &c.author_email, opts.max_commits_per_file);
walk.files.insert(to.clone(), st);
walk.dirty.remove(from);
walk.dirty.insert(to.clone());
walk.deleted.remove(to);
walk.touched_edges
.push(("TOUCHED".into(), c.sha.clone(), to.clone()));
let exists = incremental && w.has_node(from);
walk.rename(from, to, exists);
}
}
}
}
work.touched.extend(walk.dirty.iter().cloned());
for key in walk.deleted.iter().chain(walk.alias.keys()) {
if !walk.files.contains_key(key) {
work.stale.insert(key.clone());
}
}
let legacy = legacy_name_counts(w, p)?;
let mut author_rows = Vec::new();
let mut author_updates = Vec::new();
for (email, seen) in &walk.authors {
let mut merged = AuthorState::default();
match (
w.node_ref(email).and_then(|n| n.prop("name_counts")),
legacy.get(email),
) {
(Some(Value::List(l)), _) => {
merged.set_name_counts(&l);
merged.absorb(seen);
}
(_, Some(whole_history)) => merged = whole_history.clone(),
_ => merged.absorb(seen),
}
let props = [
("name", Value::Str(merged.display_name())),
("name_counts", merged.name_counts_value()),
];
if w.has_node(email) {
for (field, want) in props {
if w.node_ref(email).and_then(|n| n.prop(field)).as_ref() != Some(&want) {
author_updates.push((email.clone(), field, want));
}
}
} else {
let mut row = BTreeMap::from([("id".to_string(), Value::Str(email.clone()))]);
row.extend(props.into_iter().map(|(f, v)| (f.to_string(), v)));
author_rows.push(row);
}
}
if !author_updates.is_empty() {
let mut b = w.batch();
for (key, field, value) in author_updates {
b.set_prop(&key, field, value);
}
b.commit()?;
}
let a = w.ingest_with_edges("Author", author_rows, ingest, &[])?;
report.rules_created.extend(a.rules_created);
authors.extend(walk.authors.keys().cloned());
for p in &walk.deleted {
if w.has_node(p) {
w.delete_node(p)?;
report.deleted += 1;
}
}
for (from, to) in &walk.renamed {
if !w.has_node(from) || from == to {
continue;
}
if !walk.files.contains_key(to) {
w.delete_node(from)?;
report.deleted += 1;
continue;
}
if w.has_node(to) {
w.delete_node(to)?;
report.deleted += 1;
}
w.rename_node(from, to)?;
w.set_prop(to, "id", Value::Str(to.clone()))?;
report.renamed += 1;
}
let mut new_file_rows = Vec::new();
let mut updates = Vec::new();
let mut written = 0usize;
for (path, st) in &walk.files {
if incremental && !walk.dirty.contains(path) {
continue;
}
written += 1;
let props = file_props(st, path);
if w.has_node(path) {
updates.extend(props.into_iter().map(|(k, v)| (path.clone(), k, v)));
} else {
new_file_rows.push(props.into_iter().collect::<BTreeMap<_, _>>());
}
}
if !updates.is_empty() {
let mut b = w.batch();
for (key, field, value) in updates {
b.set_prop(&key, &field, value);
}
b.commit()?;
}
let f = w.ingest_with_edges("File", new_file_rows, ingest, &[])?;
report.rules_created.extend(f.rules_created);
report.files += written;
let c = w.ingest_with_edges("Commit", walk.commit_rows, ingest, &[])?;
report.rules_created.extend(c.rules_created);
report.commits += log.len();
let touched: Vec<(String, String, String)> = walk
.touched_edges
.into_iter()
.map(|(t, sha, p)| {
let p = walk.alias.get(&p).cloned().unwrap_or(p);
(t, sha, p)
})
.filter(|(_, _, p)| walk.files.contains_key(p))
.collect();
if !touched.is_empty() {
w.ingest_with_edges("Commit", Vec::new(), ingest, &touched)?;
}
Ok(())
}
const NO_MARKER: &str = "store has no git sync marker; run ingest-git first";
#[derive(Debug, Clone, PartialEq, Eq)]
struct SyncMarker {
repo: PathBuf,
recurse: bool,
prs: bool,
structure: bool,
docs: bool,
}
impl SyncMarker {
fn opts(&self) -> IngestGitOpts {
IngestGitOpts {
repo: self.repo.clone(),
exclude: DEFAULT_EXCLUDES.iter().map(|p| (*p).to_string()).collect(),
max_commits_per_file: DEFAULT_MAX_COMMITS_PER_FILE,
recurse_submodules: self.recurse,
prs: self.prs,
structure: self.structure,
docs: self.docs,
ensure_gitignore: false,
}
}
}
fn marker_of(r: &structure::Db) -> Result<SyncMarker, CliError> {
let node = r
.node_ref(SYNC_KEY)
.ok_or_else(|| CliError(NO_MARKER.into()))?;
let flag = |name: &str| matches!(node.prop(name), Some(Value::Bool(true)));
match node.prop("repo") {
Some(Value::Str(repo)) if !repo.is_empty() => Ok(SyncMarker {
repo: PathBuf::from(repo),
recurse: flag("recurse"),
prs: flag("prs"),
structure: node.prop("structure") != Some(Value::Bool(false)),
docs: node.prop("docs") != Some(Value::Bool(false)),
}),
_ => Err(CliError(NO_MARKER.into())),
}
}
fn open_existing(db_dir: &Path) -> Result<SharedDb, CliError> {
if !db_dir.exists() {
return Err(CliError(format!(
"no database directory at {}",
db_dir.display()
)));
}
Ok(SharedDb::open(db_dir)?)
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct SyncReport {
pub git: IngestGitReport,
pub structure: crate::structure::StructureReport,
pub dirty_refreshed: usize,
}
fn dirty_paths(repo: &Path, exclude: &[String]) -> Result<Vec<String>, CliError> {
const LISTS: [&[&str]; 2] = [
&["diff", "--name-only", "-z", "HEAD"],
&["ls-files", "--others", "--exclude-standard", "-z"],
];
let mut out = BTreeSet::new();
for args in LISTS {
let o = git_output(repo, args)?;
if !o.status.success() {
continue;
}
for path in String::from_utf8_lossy(&o.stdout).split('\0') {
if path.is_empty() || excluded(path, exclude) {
continue;
}
out.insert(path.to_string());
}
}
Ok(out.into_iter().collect())
}
pub fn run_sync(db_dir: &Path) -> Result<SyncReport, CliError> {
let db = open_existing(db_dir)?;
let marker = marker_of(&db.read())?;
let opts = marker.opts();
let mut report = SyncReport {
git: run_ingest_git(db_dir, &opts)?,
..Default::default()
};
if !opts.structure {
return Ok(report);
}
let repo = canonical(&marker.repo);
let paths = dirty_paths(&repo, &opts.exclude)?;
report.dirty_refreshed = paths.len();
if paths.is_empty() {
return Ok(report);
}
let mut w = db.write_with_wait(WRITE_LOCK_WAIT).map_err(|e| match e {
GraphError::Busy { .. } => CliError(BUSY_MESSAGE.to_string()),
other => CliError(other.to_string()),
})?;
report.structure = structure::refresh_files(&mut w, &repo, "", &paths, opts.docs)?;
drop(w);
snapshot_if_due(&db, db_dir, false);
Ok(report)
}
pub fn format_touch(r: &structure::StructureReport) -> String {
format!(
"touch: {} file(s), {} symbol(s), {} import(s), {} call(s), {} mention(s)\n",
r.files_scanned, r.symbols, r.imports, r.calls, r.mentions
)
}
#[must_use]
pub fn format_sync_json(r: &SyncReport) -> String {
let g = &r.git;
let s = &r.structure;
let gs = &g.structure;
let value = serde_json::json!({
"text": format_sync(r),
"commits": g.commits,
"files": g.files,
"authors": g.authors,
"renamed": g.renamed,
"deleted": g.deleted,
"incremental": g.incremental,
"submodules": g.submodules,
"prs": g.prs,
"rules_created": g.rules_created,
"scanned": {
"files": gs.files_scanned,
"symbols": gs.symbols,
"imports": gs.imports,
"calls": gs.calls,
"mentions": gs.mentions,
},
"dirty_refreshed": r.dirty_refreshed,
"dirty": {
"files": s.files_scanned,
"symbols": s.symbols,
"imports": s.imports,
"calls": s.calls,
"mentions": s.mentions,
},
});
format!("{value}\n")
}
pub fn format_sync(r: &SyncReport) -> String {
let mut out = format_ingest_git(&r.git);
let s = &r.structure;
out.push_str(&format!(
" dirty {} path(s): scanned {}, {} symbol(s), {} import(s), {} call(s)\n",
r.dirty_refreshed, s.files_scanned, s.symbols, s.imports, s.calls
));
out
}
pub fn run_touch(
db_dir: &Path,
files: &[PathBuf],
hook_stdin: Option<&str>,
) -> Result<structure::StructureReport, CliError> {
let named: Vec<PathBuf> = if files.is_empty() {
hook_stdin.map(paths_from_payload).unwrap_or_default()
} else {
files.to_vec()
};
if named.is_empty() {
return Ok(structure::StructureReport::default());
}
let db = open_existing(db_dir)?;
let marker = marker_of(&db.read())?;
if !marker.structure {
return Ok(structure::StructureReport::default());
}
let repo = canonical(&marker.repo);
let exclude: Vec<String> = DEFAULT_EXCLUDES.iter().map(|p| (*p).to_string()).collect();
let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
let mut paths: BTreeSet<String> = BTreeSet::new();
for path in &named {
let Some(rel) = repo_relative(&repo, &cwd, path) else {
continue;
};
if !excluded(&rel, &exclude) {
paths.insert(rel);
}
}
if paths.is_empty() {
return Ok(structure::StructureReport::default());
}
let paths: Vec<String> = paths.into_iter().collect();
let mut w = db.write_with_wait(WRITE_LOCK_WAIT).map_err(|e| match e {
GraphError::Busy { .. } => CliError(BUSY_MESSAGE.to_string()),
other => CliError(other.to_string()),
})?;
structure::refresh_files(&mut w, &repo, "", &paths, marker.docs)
}
fn paths_from_payload(raw: &str) -> Vec<PathBuf> {
let Ok(v) = serde_json::from_str::<serde_json::Value>(raw) else {
return Vec::new();
};
let input = &v["tool_input"];
let mut out = Vec::new();
let mut push = |value: &serde_json::Value| {
if let Some(s) = value.as_str().map(str::trim).filter(|s| !s.is_empty()) {
out.push(PathBuf::from(s));
}
};
push(&input["file_path"]);
if let Some(edits) = input["edits"].as_array() {
for e in edits {
push(&e["file_path"]);
}
}
out
}
fn repo_relative(repo: &Path, cwd: &Path, path: &Path) -> Option<String> {
let absolute = if path.is_absolute() {
path.to_path_buf()
} else {
cwd.join(path)
};
let resolved = resolve_symlinks(&absolute);
let rel = resolved.strip_prefix(repo).ok()?;
let key: Vec<String> = rel
.components()
.map(|c| c.as_os_str().to_string_lossy().into_owned())
.collect();
let key = key.join("/");
(!key.is_empty()).then_some(key)
}
fn resolve_symlinks(p: &Path) -> PathBuf {
if let Ok(resolved) = std::fs::canonicalize(p) {
return resolved;
}
match (p.parent(), p.file_name()) {
(Some(dir), Some(name)) => std::fs::canonicalize(dir)
.map(|d| d.join(name))
.unwrap_or_else(|_| p.to_path_buf()),
_ => p.to_path_buf(),
}
}
pub fn format_ingest_git(r: &IngestGitReport) -> String {
let mut out = format!(
"ingest-git: {} commit(s), {} file(s), {} author(s){}\n",
r.commits,
r.files,
r.authors,
if r.incremental { " (incremental)" } else { "" }
);
if r.renamed + r.deleted > 0 {
out.push_str(&format!(" renamed {} deleted {}\n", r.renamed, r.deleted));
}
if r.submodules + r.prs > 0 {
out.push_str(&format!(
" submodules {} pull requests {}\n",
r.submodules, r.prs
));
}
let s = &r.structure;
if s.files_scanned > 0 {
out.push_str(&format!(
" scanned {} file(s): {} symbol(s), {} import(s), {} call(s), {} mention(s)\n",
s.files_scanned, s.symbols, s.imports, s.calls, s.mentions
));
if s.skipped_large + s.symbols_capped > 0 {
out.push_str(&format!(
" hash-only {} symbol cap hit on {}\n",
s.skipped_large, s.symbols_capped
));
}
}
if r.gitignore_added {
out.push_str(" added the database directory to .gitignore\n");
}
if !r.rules_created.is_empty() {
out.push_str(&format!(" rules: {}\n", r.rules_created.join(", ")));
}
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn exclude_matches_prefix_extension_and_substring() {
let pats = vec![
"target/".to_string(),
"*.lock".into(),
"node_modules".into(),
];
assert!(excluded("target/debug/foo.rs", &pats));
assert!(
!excluded("targeted/foo.rs", &pats),
"prefix needs the slash"
);
assert!(excluded("Cargo.lock", &pats));
assert!(!excluded("Cargo.toml", &pats));
assert!(excluded("ui/node_modules/x/y.js", &pats));
assert!(!excluded("src/lib.rs", &pats));
assert!(!excluded("anything", &[]));
}
#[test]
fn a_compound_suffix_pattern_matches() {
let defaults: Vec<String> = DEFAULT_EXCLUDES.iter().map(|p| (*p).to_string()).collect();
assert!(excluded("ui/build/bundle.min.js", &defaults));
assert!(excluded("bundle.min.js", &defaults));
assert!(
!excluded("ui/src/app.js", &defaults),
"an ordinary source file is not a bundle"
);
assert!(!excluded("ui/src/minify.js", &defaults));
assert!(excluded("Cargo.lock", &defaults));
assert!(!excluded(".lock", &defaults));
assert!(!excluded("src/lib.rs", &defaults));
}
#[test]
fn file_props_split_dir_and_ext() {
let mut st = FileState::default();
st.touch("sha1", "a@x.test", 200);
let m: BTreeMap<_, _> = file_props(&st, "src/a/b.rs").into_iter().collect();
assert_eq!(m["dir"], Value::Str("src/a".into()));
assert_eq!(m["ext"], Value::Str("rs".into()));
assert_eq!(m["n_commits"], Value::Int(1));
assert_eq!(m["id"], Value::Str("src/a/b.rs".into()));
assert_eq!(m["top_author_id"], Value::Str("a@x.test".into()));
assert_eq!(
m["author_counts"],
Value::List(vec![Value::Str("a@x.test\t1".into())])
);
let m: BTreeMap<_, _> = file_props(&FileState::default(), "README")
.into_iter()
.collect();
assert_eq!(m["dir"], Value::Str(String::new()));
assert_eq!(m["ext"], Value::Str(String::new()));
}
#[test]
fn author_counts_round_trip_preserves_the_distribution() {
let mut st = FileState::default();
for _ in 0..3 {
st.touch("s", "alice@x.test", 200);
}
for _ in 0..4 {
st.touch("s", "bob@x.test", 200);
}
let Value::List(encoded) = st.author_counts_value() else {
panic!("author_counts must be a list");
};
assert_eq!(
encoded,
vec![
Value::Str("alice@x.test\t3".into()),
Value::Str("bob@x.test\t4".into()),
],
"email order, so the prop is stable across runs"
);
let mut reloaded = FileState::default();
reloaded.set_author_counts(&encoded);
assert_eq!(reloaded.by_author, st.by_author);
assert_eq!(reloaded.top_author(), "bob@x.test");
}
#[test]
fn author_counts_skips_entries_it_cannot_parse() {
let mut st = FileState::default();
st.set_author_counts(&[
Value::Str("alice@x.test\t2".into()),
Value::Str("no-tab-here".into()),
Value::Str("bob@x.test\tnotanumber".into()),
Value::Str("\t5".into()),
Value::Int(7),
]);
assert_eq!(st.by_author, BTreeMap::from([("alice@x.test".into(), 2)]));
}
#[test]
fn commits_list_is_capped_and_top_author_is_deterministic() {
let mut st = FileState::default();
for i in 0..5 {
st.touch(&format!("sha{i}"), "b@x.test", 3);
}
st.touch("shaX", "a@x.test", 3);
assert_eq!(st.commits, vec!["sha3", "sha4", "shaX"]);
assert_eq!(st.n_commits, 6);
assert_eq!(st.top_author(), "b@x.test");
let m: BTreeMap<_, _> = file_props(&st, "src/hot.rs").into_iter().collect();
assert_eq!(m["n_commits"], Value::Int(6));
let Value::List(commits) = &m["commits"] else {
panic!("commits must be a list");
};
assert_eq!(commits.len(), 3, "the list is still capped at 3");
assert!(
matches!(m["n_commits"], Value::Int(n) if n as usize > commits.len()),
"n_commits must exceed the capped list once the cap is passed"
);
assert_eq!(
m["author_counts"],
Value::List(vec![
Value::Str("a@x.test\t1".into()),
Value::Str("b@x.test\t5".into()),
]),
"the per-author counts sum to n_commits, not to the capped list"
);
let mut tie = FileState::default();
tie.touch("s", "b@x.test", 10);
tie.touch("s", "a@x.test", 10);
assert_eq!(tie.top_author(), "a@x.test", "ties break on smallest email");
}
}