use crate::CliError;
use core_api::{Direction, IngestOptions, Predicate, ResultSet, RuleDef, SharedDb, Value};
use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::process::Command;
pub const DEFAULT_MAX_COMMITS_PER_FILE: usize = 200;
const CO_CHANGE_MIN: f64 = 0.25;
const SYNC_KEY: &str = "__mushroomdb_git_sync__";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngestGitOpts {
pub repo: PathBuf,
pub exclude: Vec<String>,
pub max_commits_per_file: usize,
}
#[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>,
}
#[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 {
patterns.iter().any(|p| {
if let Some(prefix) = p.strip_suffix('/') {
path.starts_with(&format!("{prefix}/"))
} else if let Some(ext) = p.strip_prefix("*.") {
path.contains('.') && path.rsplit('.').next() == Some(ext)
} else {
path.contains(p.as_str())
}
})
}
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))
}
const AUTHOR_COUNT_SEP: char = '\t';
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;
}
}
}
}
const FILE_STATE_QUERY: &str = "MATCH (f:File) 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) -> 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,
};
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
}
#[derive(Default)]
struct Walk {
files: BTreeMap<String, FileState>,
authors: BTreeMap<String, String>,
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());
}
}
pub fn run_ingest_git(db_dir: &Path, opts: &IngestGitOpts) -> Result<IngestGitReport, CliError> {
let db = SharedDb::open(db_dir)?;
let since: Option<String> = {
let r = db.read();
r.node_ref(SYNC_KEY)
.and_then(|n| n.prop("sha"))
.and_then(|v| match v {
Value::Str(s) if !s.is_empty() => Some(s),
_ => None,
})
};
let incremental = since.is_some();
let mut report = IngestGitReport {
incremental,
..Default::default()
};
let Some(head) = head_sha(&opts.repo)? else {
return Ok(report); };
let log = read_log(&opts.repo, since.as_deref(), &head)?;
if log.is_empty() {
return Ok(report);
}
let mut w = db.write();
let ingest = IngestOptions::default();
let mut walk = Walk {
files: if incremental {
file_state_from(&w.query(FILE_STATE_QUERY, &BTreeMap::new())?)
} else {
BTreeMap::new()
},
..Default::default()
};
for c in &log {
walk.authors
.entry(c.author_email.clone())
.or_insert_with(|| c.author_name.clone());
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);
}
}
}
}
let author_rows: Vec<BTreeMap<String, Value>> = walk
.authors
.iter()
.filter(|(email, _)| !w.has_node(email))
.map(|(email, name)| {
BTreeMap::from([
("id".to_string(), Value::Str(email.clone())),
("name".to_string(), Value::Str(name.clone())),
])
})
.collect();
let a = w.ingest_with_edges("Author", author_rows, &ingest, &[])?;
report.rules_created.extend(a.rules_created);
report.authors = walk.authors.len();
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();
for (path, st) in &walk.files {
if incremental && !walk.dirty.contains(path) {
continue;
}
let props = file_props(st, path);
if w.has_node(path) {
for (k, v) in props {
w.set_prop(path, &k, v)?;
}
} else {
new_file_rows.push(props.into_iter().collect::<BTreeMap<_, _>>());
}
}
let f = w.ingest_with_edges("File", new_file_rows, &ingest, &[])?;
report.rules_created.extend(f.rules_created);
report.files = walk.files.len();
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)?;
}
if !incremental {
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,
})?;
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),
})?;
report
.rules_created
.extend(["co_changed".to_string(), "knows".to_string()]);
for (l, field) in [("File", "path"), ("Commit", "message"), ("Author", "name")] {
w.enable_fulltext(l, field)?;
}
}
if w.has_node(SYNC_KEY) {
w.set_prop(SYNC_KEY, "sha", Value::Str(head))?;
} else {
w.insert_node(
"GitSync",
SYNC_KEY,
vec![
("id".into(), Value::Str(SYNC_KEY.into())),
("sha".into(), Value::Str(head)),
],
)?;
}
Ok(report)
}
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.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 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");
}
}