cli/ingest_git.rs
1//! `mushroomdb ingest-git <db> <repo>`: build and maintain a graph of a git
2//! repository. First run ingests the whole history; later runs apply only the
3//! commits after the recorded `GitSync` head, so deletes and renames retract
4//! or move derived edges instead of leaving them stale.
5//!
6//! Graph shape:
7//!
8//! | Label | key (`id`) | props |
9//! |---|---|---|
10//! | `Author` | mailmap-resolved email | `name`, `name_counts` |
11//! | `Commit` | full sha | `message`, `ts`, `author_id`, `pr_id` (with `--prs`) |
12//! | `File` | path | `path`, `dir`, `ext`, `commits`, `n_commits`, `top_author_id`, `author_counts`, plus the working-tree props in [`structure`](crate::structure) |
13//! | `Symbol` | `"<path>#<qualified name>"` | see [`structure`](crate::structure) |
14//! | `PR` | `"pr:<number>"` | `number`, `title`, `url`, `merged_at`, `author_login` |
15//! | `GitSync` | `"__mushroomdb_git_sync__"` | `sha`, `synced_at`, `repo`, `recurse`, `prs`, `structure`, `docs` |
16//!
17//! Edges: user `TOUCHED` Commit→File and `MERGED_AS` PR→Commit, auto-FK
18//! `AUTHOR` Commit→Author, `TOP_AUTHOR` File→Author and `PR` Commit→PR,
19//! rule-derived `CO_CHANGED` File→File and `KNOWS` Author→File — and, from the
20//! working-tree pass, `DEFINES`, `IMPORTS`, `CALLS` and `MENTIONS`.
21//!
22//! With `--recurse-submodules` each initialised submodule is walked as its own
23//! *unit*: its file keys carry the submodule's path in the parent, and it
24//! resumes from its own `GitSync` marker.
25use crate::structure;
26use crate::CliError;
27use core_api::{
28 default_max_edges, Direction, GraphError, IngestOptions, Predicate, ResultSet, RuleDef,
29 SharedDb, Value, WriteGuard, WRITE_LOCK_WAIT,
30};
31use std::collections::{BTreeMap, BTreeSet};
32use std::path::{Path, PathBuf};
33use std::process::Command;
34
35/// Cap on the `commits` list stored per file. Bounds both node size and the
36/// cost of the jaccard overlap the `co_changed` rule runs over that list.
37pub const DEFAULT_MAX_COMMITS_PER_FILE: usize = 200;
38
39/// The paths a repository carries that are not its source: build output,
40/// vendored dependencies, generated bundles, and lockfiles nobody reads.
41///
42/// Applied when the user names no `--exclude` pattern of their own, which keeps
43/// them out of the history graph *and* out of the working-tree pass.
44pub const DEFAULT_EXCLUDES: [&str; 6] = [
45 "target/",
46 "node_modules/",
47 "dist/",
48 ".git/",
49 "*.lock",
50 "*.min.js",
51];
52
53/// Minimum jaccard overlap of two files' `commits` lists for `CO_CHANGED`.
54const CO_CHANGE_MIN: f64 = 0.25;
55
56/// WAL bytes past the last snapshot that make writing a new one worthwhile.
57///
58/// Every open reads `wal.bin` whole and replays it frame by frame, so the tail
59/// is paid again on every hook, every MCP start and every CLI call — while a
60/// snapshot is paid once by the run that writes it. Measured on this
61/// repository, an 8.1 MB tail costs 156 ms per open against ~440 ms to
62/// snapshot it away, so a tail this size pays its snapshot back within three
63/// opens. Below the threshold the replay is cheap enough that snapshotting on
64/// every incremental run would cost more than it saves.
65pub const SNAPSHOT_WAL_BYTES: u64 = 4 * 1024 * 1024;
66
67/// Key of the singleton `GitSync` node holding the last ingested sha.
68///
69/// Node keys are a single namespace shared with `File` keys, which are repo
70/// paths — so this cannot be `"HEAD"`. A repository with a file named `HEAD`
71/// (git's own `.git/HEAD` aside, plenty of projects ship one) would otherwise
72/// have the sha written onto its `File` node, leaving no sync marker and
73/// forcing a full re-ingest on every run.
74pub(crate) const SYNC_KEY: &str = "__mushroomdb_git_sync__";
75
76/// What the CLI prints, and exits 3 on, when another process holds the store's
77/// cross-process write lock. Nothing was written, so retrying is always safe.
78pub const BUSY_MESSAGE: &str = "another mushroomdb process is writing; retry";
79
80/// Name of the `Commit.pr_id` foreign-key rule. Identical to the name the
81/// zero-config FK inference would choose, so the two never both create it.
82const PR_FK_RULE: &str = "auto_fk_commit_pr_id";
83
84#[derive(Debug, Clone, PartialEq, Eq)]
85pub struct IngestGitOpts {
86 pub repo: PathBuf,
87 /// Paths to skip. Pattern ending in `/` = path prefix, pattern starting
88 /// with `*.` = extension, otherwise = substring of the path.
89 pub exclude: Vec<String>,
90 pub max_commits_per_file: usize,
91 /// Walk every initialised submodule as its own unit.
92 pub recurse_submodules: bool,
93 /// Ask `gh` for merged pull requests and link them to their commits.
94 pub prs: bool,
95 /// Read the working tree: content hashes, `Symbol` nodes, imports and
96 /// calls. Off with `--no-structure`. Recorded on the `GitSync` node.
97 pub structure: bool,
98 /// Index Markdown bodies, headings and mentions. Off with `--no-docs`, and
99 /// inert without `structure`. Recorded on the `GitSync` node.
100 pub docs: bool,
101 /// Add the database directory to the repository's `.gitignore`.
102 pub ensure_gitignore: bool,
103}
104
105#[derive(Debug, Default, Clone, PartialEq, Eq)]
106pub struct IngestGitReport {
107 pub commits: usize,
108 pub files: usize,
109 pub authors: usize,
110 pub renamed: usize,
111 /// Nodes removed because the window deleted the file they named.
112 pub deleted: usize,
113 /// Nodes removed because a rename destination did not survive the window:
114 /// the file was renamed, and the path it moved to was deleted later in the
115 /// same walk. The node goes with it, but no path this window deleted ever
116 /// named that node, so counting it as a delete overstated how much the
117 /// window removed.
118 ///
119 /// A rename whose destination is *excluded* is not this — it is classified
120 /// upstream as a plain delete of the source path (`walk.deleted`), and
121 /// lands in [`deleted`](Self::deleted).
122 pub evicted: usize,
123 pub incremental: bool,
124 pub rules_created: Vec<String>,
125 /// Submodules walked as their own units.
126 pub submodules: usize,
127 /// Pull requests inserted by this run.
128 pub prs: usize,
129 /// Whether this run appended the database directory to `.gitignore`.
130 pub gitignore_added: bool,
131 /// What the working-tree pass saw. All zeros with `--no-structure`.
132 pub structure: crate::structure::StructureReport,
133}
134
135/// Paths the commit walk left for the working-tree pass to look at.
136#[derive(Default)]
137struct StructureWork {
138 /// Files added, modified, or renamed *into* this window's paths.
139 touched: BTreeSet<String>,
140 /// Keys that no longer name a file: renamed away, or deleted. Any file
141 /// whose `imports` or `mentions` list still holds one of these has to be
142 /// extracted again, or the edge it derived stays behind.
143 stale: BTreeSet<String>,
144}
145
146/// One git working tree walked by a run: the repository itself, or one of its
147/// submodules.
148///
149/// A submodule's paths are keys under `prefix` (its path in the parent), and it
150/// carries its own sync marker, so the two histories advance independently.
151#[derive(Debug, Clone)]
152struct RepoUnit {
153 path: PathBuf,
154 /// `""` for the repository itself, `"<displaypath>/"` for a submodule.
155 prefix: String,
156 sync_key: String,
157}
158
159#[derive(Debug)]
160enum Change {
161 Added(String),
162 Modified(String),
163 Deleted(String),
164 Renamed { from: String, to: String },
165}
166
167#[derive(Debug)]
168struct GitCommit {
169 sha: String,
170 author_name: String,
171 author_email: String,
172 ts: i64,
173 subject: String,
174 changes: Vec<Change>,
175}
176
177/// Whether `path` matches any of `patterns`. Documented in `docs/site/ingest-git.md`.
178///
179/// A `foo/` pattern is a *directory prefix*. A `*.` pattern is a **file-name
180/// suffix**, not a single extension: `*.min.js` matches `ui/bundle.min.js` the
181/// same way `*.lock` matches `Cargo.lock`. Matching only the last dot segment
182/// would leave every compound suffix inert, and a compound suffix is exactly
183/// how generated files announce themselves. Anything else is a substring.
184#[must_use]
185pub fn path_excluded(path: &str, patterns: &[String]) -> bool {
186 patterns.iter().any(|p| {
187 if let Some(prefix) = p.strip_suffix('/') {
188 path.starts_with(&format!("{prefix}/"))
189 } else if let Some(suffix) = p.strip_prefix('*').filter(|s| s.starts_with('.')) {
190 // The suffix must follow something, so `*.lock` does not claim a
191 // path that is nothing but the suffix itself.
192 path.len() > suffix.len() && path.ends_with(suffix)
193 } else {
194 path.contains(p.as_str())
195 }
196 })
197}
198
199fn git_output(repo: &Path, args: &[&str]) -> Result<std::process::Output, CliError> {
200 Command::new("git")
201 .arg("-C")
202 .arg(repo)
203 .args(args)
204 .output()
205 .map_err(|e| CliError(format!("cannot run git in {}: {e}", repo.display())))
206}
207
208/// `git log --reverse --name-status -M --format=<RS>%H<US>%aN<US>%aE<US>%at<US>%s <range>`
209///
210/// Returns oldest commit first. The walk ends at `head` — never at the symbolic
211/// `HEAD` — so the range is pinned to the same sha the caller will record as the
212/// sync marker. See [`head_sha`] for why that matters.
213///
214/// `%aN` and `%aE` are the mailmap-resolved name and address, so a repository
215/// with a `.mailmap` reports one identity for a contributor who has committed
216/// under several addresses. Without one they are exactly `%an`/`%ae`.
217fn read_log(repo: &Path, since: Option<&str>, head: &str) -> Result<Vec<GitCommit>, CliError> {
218 if let Some(s) = since {
219 let spec = format!("{s}^{{commit}}");
220 if !git_output(repo, &["cat-file", "-e", &spec])?
221 .status
222 .success()
223 {
224 return Err(CliError(format!(
225 "recorded sync head {s} is not in {} (history rewritten?); \
226 ingest into a fresh database directory",
227 repo.display()
228 )));
229 }
230 }
231 let mut cmd = Command::new("git");
232 cmd.arg("-C").arg(repo).args([
233 // Without this git renders any non-ASCII byte in a path as an octal
234 // escape, so `src/café.rs` would be stored under a mangled key that no
235 // later run matches. Paths containing a tab or newline stay quoted and
236 // escaped either way — git has to, or they would break the format below.
237 "-c",
238 "core.quotePath=false",
239 "log",
240 "--reverse",
241 "--name-status",
242 "-M",
243 "--no-color",
244 // Record separator \x1e between commits, unit separator \x1f between
245 // header fields. A commit subject containing either byte splits its own
246 // record: the message is truncated at the first \x1f, and a \x1e drops
247 // the remainder of that commit's header. Accepted — the parse degrades
248 // to a skipped or shortened message, never a panic or a wrong sha.
249 "--format=%x1e%H%x1f%aN%x1f%aE%x1f%at%x1f%s",
250 ]);
251 match since {
252 Some(s) => cmd.arg(format!("{s}..{head}")),
253 None => cmd.arg(head),
254 };
255 let out = cmd
256 .output()
257 .map_err(|e| CliError(format!("cannot run git: {e}")))?;
258 if !out.status.success() {
259 return Err(CliError(format!(
260 "git log failed: {}",
261 String::from_utf8_lossy(&out.stderr).trim()
262 )));
263 }
264 let text = String::from_utf8_lossy(&out.stdout);
265 let mut commits = Vec::new();
266 for block in text.split('\x1e').filter(|b| !b.trim().is_empty()) {
267 let mut lines = block.lines();
268 let header = lines.next().unwrap_or("");
269 let f: Vec<&str> = header.split('\x1f').collect();
270 if f.len() < 5 {
271 continue;
272 }
273 let mut changes = Vec::new();
274 for l in lines {
275 let cols: Vec<&str> = l.split('\t').collect();
276 match cols.as_slice() {
277 [s, p] if s.starts_with('A') => changes.push(Change::Added((*p).to_string())),
278 [s, p] if s.starts_with('M') || s.starts_with('T') => {
279 changes.push(Change::Modified((*p).to_string()))
280 }
281 [s, p] if s.starts_with('D') => changes.push(Change::Deleted((*p).to_string())),
282 [s, from, to] if s.starts_with('R') => changes.push(Change::Renamed {
283 from: (*from).to_string(),
284 to: (*to).to_string(),
285 }),
286 // A copy is a brand-new path with no prior history here.
287 [s, _from, to] if s.starts_with('C') => {
288 changes.push(Change::Added((*to).to_string()))
289 }
290 _ => {}
291 }
292 }
293 commits.push(GitCommit {
294 sha: f[0].into(),
295 author_name: f[1].into(),
296 author_email: f[2].into(),
297 ts: f[3].parse().unwrap_or(0),
298 subject: f[4].into(),
299 changes,
300 });
301 }
302 Ok(commits)
303}
304
305/// Resolve `HEAD` to a concrete sha. `Ok(None)` means the repository has no
306/// commits yet; a path that is not a repository at all is an error.
307///
308/// Called **before** [`read_log`], and the resulting sha is both the end of the
309/// walk and the recorded sync marker. Resolving it afterwards instead would open
310/// a window: a commit landing between the walk and the `rev-parse` would push
311/// the marker past a commit that was never ingested, and every later run would
312/// skip it silently. Pinning both to one sha closes that window — a commit that
313/// lands mid-run is simply outside this range and gets picked up next time.
314///
315/// Asking git also beats taking `log.last()`. The two agree in practice, since a
316/// reachability walk emits its tip first and reversing puts it last, but that is
317/// a property of the traversal and of this module's parser rather than something
318/// the resume marker should rest on.
319fn head_sha(repo: &Path) -> Result<Option<String>, CliError> {
320 let out = git_output(repo, &["rev-parse", "--verify", "-q", "HEAD^{commit}"])?;
321 if !out.status.success() {
322 // No commits yet, or not a repository at all — only the latter is an error.
323 if !git_output(repo, &["rev-parse", "--git-dir"])?
324 .status
325 .success()
326 {
327 return Err(CliError(format!(
328 "not a git repository: {}",
329 repo.display()
330 )));
331 }
332 return Ok(None);
333 }
334 let sha = String::from_utf8_lossy(&out.stdout).trim().to_string();
335 if sha.is_empty() {
336 return Err(CliError(format!(
337 "git could not resolve HEAD in {}",
338 repo.display()
339 )));
340 }
341 Ok(Some(sha))
342}
343
344/// Absolute, symlink-free form of `p`, falling back to `p` when it cannot be
345/// resolved (a path that does not exist yet, or a permission error). The result
346/// is what `GitSync.repo` records, so a later run can find the repository again
347/// from any working directory.
348fn canonical(p: &Path) -> PathBuf {
349 std::fs::canonicalize(p).unwrap_or_else(|_| p.to_path_buf())
350}
351
352/// The submodule paths recorded in a unit's `.gitmodules`, relative to that
353/// unit.
354///
355/// These are the paths git reports as ordinary changes in `--name-status`
356/// output while being gitlinks — a commit pointer, not a file. They get no
357/// `File` node whether or not the submodule is walked, and the list is read
358/// from configuration so it is the same for an uninitialised submodule.
359fn gitlink_paths(unit: &Path) -> BTreeSet<String> {
360 let out = Command::new("git")
361 .arg("config")
362 .arg("--file")
363 .arg(unit.join(".gitmodules"))
364 .args(["--get-regexp", r"^submodule\..*\.path$"])
365 .output();
366 let Ok(out) = out else { return BTreeSet::new() };
367 if !out.status.success() {
368 return BTreeSet::new(); // no .gitmodules, so no submodules
369 }
370 String::from_utf8_lossy(&out.stdout)
371 .lines()
372 .filter_map(|l| l.split_once(' '))
373 .map(|(_, path)| path.trim().to_string())
374 .filter(|p| !p.is_empty())
375 .collect()
376}
377
378/// The repository itself, then one unit per initialised submodule when
379/// `recurse` is set.
380///
381/// `git submodule foreach` visits only initialised submodules and says nothing
382/// about the rest, which is the behaviour wanted here: a submodule that was
383/// never checked out has no working tree to walk. `$displaypath` is relative to
384/// the top-level repository, so it is exactly the key prefix its files need.
385fn repo_units(repo: &Path, recurse: bool) -> Result<Vec<RepoUnit>, CliError> {
386 let mut units = vec![RepoUnit {
387 path: repo.to_path_buf(),
388 prefix: String::new(),
389 sync_key: SYNC_KEY.to_string(),
390 }];
391 if !recurse {
392 return Ok(units);
393 }
394 let out = git_output(
395 repo,
396 &[
397 "submodule",
398 "foreach",
399 "--quiet",
400 "--recursive",
401 "echo \"$displaypath\"",
402 ],
403 )?;
404 if !out.status.success() {
405 return Ok(units);
406 }
407 let mut paths: Vec<String> = String::from_utf8_lossy(&out.stdout)
408 .lines()
409 .map(|l| l.trim_end_matches('/').trim().to_string())
410 .filter(|l| !l.is_empty())
411 .collect();
412 paths.sort();
413 paths.dedup();
414 for dp in paths {
415 let path = repo.join(&dp);
416 // `foreach` already skips the uninitialised, but a stale entry or a
417 // removed checkout would otherwise fail the whole run.
418 if !git_output(&path, &["rev-parse", "--git-dir"])?
419 .status
420 .success()
421 {
422 continue;
423 }
424 units.push(RepoUnit {
425 prefix: format!("{dp}/"),
426 sync_key: format!("{SYNC_KEY}:{dp}"),
427 path,
428 });
429 }
430 Ok(units)
431}
432
433/// Rewrite one unit's changes into repository-wide keys: drop gitlinks, then
434/// prepend the unit's prefix so a submodule's `src/lib.rs` is stored under
435/// `vendor/lib/src/lib.rs`.
436///
437/// Done before anything else looks at the log, so exclusion patterns, the
438/// stored keys and the `File` state query all speak the same path.
439fn localise(log: &mut [GitCommit], prefix: &str, gitlinks: &BTreeSet<String>) {
440 let keep = |p: &String| !gitlinks.contains(p.as_str());
441 for c in log.iter_mut() {
442 c.changes.retain(|ch| match ch {
443 Change::Added(p) | Change::Modified(p) | Change::Deleted(p) => keep(p),
444 Change::Renamed { from, to } => keep(from) && keep(to),
445 });
446 if prefix.is_empty() {
447 continue;
448 }
449 for ch in c.changes.iter_mut() {
450 match ch {
451 Change::Added(p) | Change::Modified(p) | Change::Deleted(p) => {
452 *p = format!("{prefix}{p}")
453 }
454 Change::Renamed { from, to } => {
455 *from = format!("{prefix}{from}");
456 *to = format!("{prefix}{to}");
457 }
458 }
459 }
460 }
461}
462
463/// Append `<db dir>/` to the repository's `.gitignore` unless it is already
464/// listed, creating the file if it does not exist. Returns whether it was
465/// written.
466///
467/// A database kept outside the repository is left alone: the repository has no
468/// business ignoring a path it does not contain.
469fn ensure_gitignore(repo: &Path, db_dir: &Path) -> Result<bool, CliError> {
470 let Ok(rel) = canonical(db_dir)
471 .strip_prefix(repo)
472 .map(|p| p.to_path_buf())
473 else {
474 return Ok(false);
475 };
476 if rel.as_os_str().is_empty() {
477 return Ok(false);
478 }
479 let line = format!("{}/", rel.to_string_lossy().replace('\\', "/"));
480 let path = repo.join(".gitignore");
481 let current = match std::fs::read_to_string(&path) {
482 Ok(s) => s,
483 Err(e) if e.kind() == std::io::ErrorKind::NotFound => String::new(),
484 Err(e) => return Err(CliError(format!("cannot read {}: {e}", path.display()))),
485 };
486 let bare = line.trim_end_matches('/');
487 if current
488 .lines()
489 .map(|l| l.trim())
490 .any(|l| l == line || l == bare || l == format!("/{line}") || l == format!("/{bare}"))
491 {
492 return Ok(false);
493 }
494 let mut next = current;
495 if !next.is_empty() && !next.ends_with('\n') {
496 next.push('\n');
497 }
498 next.push_str(&line);
499 next.push('\n');
500 std::fs::write(&path, next)
501 .map_err(|e| CliError(format!("cannot write {}: {e}", path.display())))?;
502 Ok(true)
503}
504
505/// Separator inside an `author_counts` or `name_counts` entry. Neither an email
506/// address nor a git author name can contain a tab, so `<text>\tcount`
507/// round-trips unambiguously.
508const AUTHOR_COUNT_SEP: char = '\t';
509
510/// How many commits carried each spelling of one email's `%aN`, in the order
511/// the spellings were first seen.
512///
513/// One person routinely commits under more than one name — `Ada Lovelace` at
514/// work and `Ada M. Lovelace` from a laptop that was configured once and never
515/// again. Merging them on email is right, and every tool that prints an author
516/// then has to choose *which* of the names to show. Showing the first one
517/// walked means a display name decided by the oldest commit in the window,
518/// which on this repository labelled an identity with a name carried by 21% of
519/// its commits.
520#[derive(Default, Clone, Debug, PartialEq, Eq)]
521struct AuthorState {
522 /// `(name, commits)` in first-seen order, which is what breaks a tie.
523 names: Vec<(String, usize)>,
524}
525
526impl AuthorState {
527 /// Credit one more commit to `name`.
528 fn touch(&mut self, name: &str) {
529 self.add(name, 1);
530 }
531
532 /// Credit `n` more commits to `name`, appending it if it is new. Order is
533 /// preserved, so the entry that was seen first stays first.
534 fn add(&mut self, name: &str, n: usize) {
535 match self.names.iter_mut().find(|(held, _)| held == name) {
536 Some(entry) => entry.1 += n,
537 None => self.names.push((name.to_string(), n)),
538 }
539 }
540
541 /// Fold another window's counts into this one, keeping first-seen order.
542 fn absorb(&mut self, other: &AuthorState) {
543 for (name, n) in &other.names {
544 self.add(name, *n);
545 }
546 }
547
548 /// The name on the most commits. A tie goes to the name seen first, which
549 /// is why the strict `>` matters: it never unseats an equal incumbent.
550 fn display_name(&self) -> String {
551 let mut best: Option<&(String, usize)> = None;
552 for entry in &self.names {
553 if best.is_none_or(|(_, n)| entry.1 > *n) {
554 best = Some(entry);
555 }
556 }
557 best.map(|(name, _)| name.clone()).unwrap_or_default()
558 }
559
560 /// The distribution as an `Author.name_counts` prop: `"name\tcount"`
561 /// strings in first-seen order.
562 ///
563 /// Stored for the same reason `File.author_counts` is: an incremental run
564 /// sees only the new window, and the `name` prop alone cannot say how far
565 /// ahead the incumbent spelling is. Without it every sync would relabel the
566 /// identity from a handful of commits.
567 fn name_counts_value(&self) -> Value {
568 Value::List(
569 self.names
570 .iter()
571 .map(|(name, n)| Value::Str(format!("{name}{AUTHOR_COUNT_SEP}{n}")))
572 .collect(),
573 )
574 }
575
576 /// Inverse of [`AuthorState::name_counts_value`]. Entries that are not
577 /// `name<TAB>count` are skipped rather than failing the run.
578 fn set_name_counts(&mut self, list: &[Value]) {
579 for v in list {
580 let Value::Str(s) = v else { continue };
581 let Some((name, n)) = s.rsplit_once(AUTHOR_COUNT_SEP) else {
582 continue;
583 };
584 let Ok(n) = n.parse::<usize>() else { continue };
585 if !name.is_empty() {
586 self.add(name, n);
587 }
588 }
589 }
590}
591
592fn file_props(st: &FileState, path: &str) -> Vec<(String, Value)> {
593 let commits = &st.commits;
594 let dir = path
595 .rsplit_once('/')
596 .map(|(d, _)| d)
597 .unwrap_or("")
598 .to_string();
599 let ext = path
600 .rsplit_once('.')
601 .map(|(_, e)| e)
602 .unwrap_or("")
603 .to_string();
604 vec![
605 ("id".into(), Value::Str(path.into())),
606 ("path".into(), Value::Str(path.into())),
607 ("dir".into(), Value::Str(dir)),
608 ("ext".into(), Value::Str(ext)),
609 (
610 "commits".into(),
611 Value::List(commits.iter().map(|s| Value::Str(s.clone())).collect()),
612 ),
613 // The true total, which past `--max-commits-per-file` is larger than
614 // the `commits` list it is stored beside. It is also what
615 // `author_counts` sums to, so the two props agree at any history
616 // length.
617 ("n_commits".into(), Value::Int(st.n_commits as i64)),
618 ("top_author_id".into(), Value::Str(st.top_author())),
619 // Written on every touch so the next incremental run can rebuild the
620 // distribution instead of crediting the whole history to the incumbent.
621 ("author_counts".into(), st.author_counts_value()),
622 ]
623}
624
625/// In-memory per-file state accumulated while walking commits.
626#[derive(Default, Clone)]
627struct FileState {
628 commits: Vec<String>,
629 by_author: BTreeMap<String, usize>,
630 n_commits: usize,
631}
632
633impl FileState {
634 fn touch(&mut self, sha: &str, author: &str, cap: usize) {
635 self.commits.push(sha.to_string());
636 if cap > 0 && self.commits.len() > cap {
637 self.commits.remove(0);
638 }
639 self.n_commits += 1;
640 *self.by_author.entry(author.to_string()).or_default() += 1;
641 }
642
643 /// Most commits wins; ties break on the lexicographically smallest email
644 /// so the result is deterministic across runs.
645 fn top_author(&self) -> String {
646 self.by_author
647 .iter()
648 .max_by(|a, b| a.1.cmp(b.1).then(b.0.cmp(a.0)))
649 .map(|(a, _)| a.clone())
650 .unwrap_or_default()
651 }
652
653 /// The per-author distribution as a `File.author_counts` prop: a list of
654 /// `"email\tcount"` strings in email order.
655 ///
656 /// This is the state an incremental run needs and cannot recompute — the
657 /// walk only sees the new window, and `top_author_id` alone cannot say how
658 /// far ahead the incumbent is. Without it a challenger's commits reset on
659 /// every sync and ownership can never change.
660 fn author_counts_value(&self) -> Value {
661 Value::List(
662 self.by_author
663 .iter()
664 .map(|(email, n)| Value::Str(format!("{email}{AUTHOR_COUNT_SEP}{n}")))
665 .collect(),
666 )
667 }
668
669 /// Inverse of [`FileState::author_counts_value`]. Entries that are not
670 /// `email<TAB>count` are skipped rather than failing the run.
671 fn set_author_counts(&mut self, list: &[Value]) {
672 for v in list {
673 let Value::Str(s) = v else { continue };
674 let Some((email, n)) = s.rsplit_once(AUTHOR_COUNT_SEP) else {
675 continue;
676 };
677 let Ok(n) = n.parse::<usize>() else { continue };
678 if !email.is_empty() {
679 *self.by_author.entry(email.to_string()).or_default() += n;
680 }
681 }
682 }
683}
684
685/// One merged pull request as reported by `gh`.
686#[derive(Debug, Clone, PartialEq, Eq)]
687struct PullRequest {
688 number: i64,
689 title: String,
690 url: String,
691 merged_at: String,
692 author_login: String,
693 /// `mergeCommit.oid`, absent for a pull request merged some other way (or
694 /// whose merge commit has since been rewritten).
695 merge_sha: Option<String>,
696}
697
698fn pr_key(number: i64) -> String {
699 format!("pr:{number}")
700}
701
702/// The pull request number in a squash-merge subject: `\(#(\d+)\)$`.
703///
704/// Hand-rolled rather than pulled in as a dependency — the pattern is anchored
705/// at the end of the subject and made of two literals around a run of digits.
706fn subject_pr(subject: &str) -> Option<i64> {
707 let rest = subject.strip_suffix(')')?;
708 let at = rest.rfind("(#")?;
709 let digits = &rest[at + 2..];
710 if digits.is_empty() || !digits.bytes().all(|b| b.is_ascii_digit()) {
711 return None;
712 }
713 digits.parse().ok()
714}
715
716/// `gh pr list --state merged` in `repo`, or an empty list with one warning.
717///
718/// Every failure is a skip: `gh` may not be installed, the repository may have
719/// no GitHub remote, and the user may not be authenticated. None of that is a
720/// reason to fail an ingest that is otherwise complete.
721fn fetch_prs(repo: &Path) -> Vec<PullRequest> {
722 let out = Command::new("gh")
723 .current_dir(repo)
724 .args([
725 "pr",
726 "list",
727 "--state",
728 "merged",
729 "--limit",
730 "1000",
731 "--json",
732 "number,title,url,mergedAt,mergeCommit,author",
733 ])
734 .output();
735 let out = match out {
736 Ok(o) => o,
737 Err(_) => {
738 eprintln!("ingest-git: --prs skipped: gh is not on PATH");
739 return Vec::new();
740 }
741 };
742 if !out.status.success() {
743 let detail = String::from_utf8_lossy(&out.stderr)
744 .lines()
745 .find(|l| !l.trim().is_empty())
746 .unwrap_or("no detail")
747 .to_string();
748 eprintln!("ingest-git: --prs skipped: gh pr list failed: {detail}");
749 return Vec::new();
750 }
751 let parsed: serde_json::Value = match serde_json::from_slice(&out.stdout) {
752 Ok(v) => v,
753 Err(e) => {
754 eprintln!("ingest-git: --prs skipped: gh pr list output is not JSON: {e}");
755 return Vec::new();
756 }
757 };
758 let Some(items) = parsed.as_array() else {
759 eprintln!("ingest-git: --prs skipped: gh pr list did not return a list");
760 return Vec::new();
761 };
762 let mut prs: Vec<PullRequest> = items
763 .iter()
764 .filter_map(|v| {
765 let number = v.get("number")?.as_i64()?;
766 let str_at = |k: &str| {
767 v.get(k)
768 .and_then(|x| x.as_str())
769 .unwrap_or_default()
770 .to_string()
771 };
772 Some(PullRequest {
773 number,
774 title: str_at("title"),
775 url: str_at("url"),
776 merged_at: str_at("mergedAt"),
777 author_login: v
778 .get("author")
779 .and_then(|a| a.get("login"))
780 .and_then(|l| l.as_str())
781 .unwrap_or_default()
782 .to_string(),
783 merge_sha: v
784 .get("mergeCommit")
785 .and_then(|m| m.get("oid"))
786 .and_then(|o| o.as_str())
787 .filter(|s| !s.is_empty())
788 .map(str::to_string),
789 })
790 })
791 .collect();
792 // Ascending by number: the insert order, and so the node order, is the
793 // same on every run whatever order gh listed them in.
794 prs.sort_by_key(|p| p.number);
795 prs.dedup_by_key(|p| p.number);
796 prs
797}
798
799/// Insert the `PR` nodes that are new, declare the `Commit.pr_id` foreign key,
800/// and index titles for search. Runs before the commits so the FK resolves.
801fn ingest_prs(
802 w: &mut WriteGuard<'_>,
803 prs: &[PullRequest],
804 ingest: &IngestOptions,
805 report: &mut IngestGitReport,
806) -> Result<(), CliError> {
807 let rows: Vec<BTreeMap<String, Value>> = prs
808 .iter()
809 .filter(|p| !w.has_node(&pr_key(p.number)))
810 .map(|p| {
811 BTreeMap::from([
812 ("id".to_string(), Value::Str(pr_key(p.number))),
813 ("number".to_string(), Value::Int(p.number)),
814 ("title".to_string(), Value::Str(p.title.clone())),
815 ("url".to_string(), Value::Str(p.url.clone())),
816 ("merged_at".to_string(), Value::Str(p.merged_at.clone())),
817 (
818 "author_login".to_string(),
819 Value::Str(p.author_login.clone()),
820 ),
821 ])
822 })
823 .collect();
824 if !rows.is_empty() {
825 let r = w.ingest_with_edges("PR", rows, ingest, &[])?;
826 report.prs = r.inserted;
827 report.rules_created.extend(r.rules_created);
828 }
829 // Declared here rather than left to FK inference, which only fires on a
830 // batch of `Commit` rows that already carry `pr_id` — a run that links a
831 // pull request to a commit ingested earlier would otherwise leave the
832 // property with no edge behind it.
833 if !w.rules().iter().any(|r| r.name == PR_FK_RULE) {
834 let predicate = Predicate::KeyMatch {
835 field: "pr_id".into(),
836 };
837 let max_edges = Some(default_max_edges(&predicate));
838 w.create_rule(RuleDef {
839 name: PR_FK_RULE.into(),
840 src_label: "Commit".into(),
841 dst_label: "PR".into(),
842 predicate,
843 edge_type: "PR".into(),
844 weight_prop: None,
845 max_edges,
846 approximate: false,
847 via_label: None,
848 via_edge: None,
849 via_dir: None,
850 namespace: None,
851 })?;
852 report.rules_created.push(PR_FK_RULE.to_string());
853 }
854 if !w
855 .fulltext_pairs()
856 .contains(&("PR".to_string(), "title".to_string()))
857 {
858 w.enable_fulltext("PR", "title")?;
859 }
860 Ok(())
861}
862
863/// Point every commit that carries a pull request at it: the merge commit by
864/// sha, a squash merge by its `(#N)` subject. Runs after the commits are in the
865/// graph, so a pull request merged before the last sync is linked too.
866fn link_prs(
867 w: &mut WriteGuard<'_>,
868 prs: &[PullRequest],
869 ingest: &IngestOptions,
870) -> Result<(), CliError> {
871 let by_sha: BTreeMap<&str, i64> = prs
872 .iter()
873 .filter_map(|p| p.merge_sha.as_deref().map(|s| (s, p.number)))
874 .collect();
875 let known: BTreeSet<i64> = prs.iter().map(|p| p.number).collect();
876
877 let rs = w.query(
878 "MATCH (c:Commit) RETURN c.id AS id, c.message AS message, c.pr_id AS pr_id",
879 &BTreeMap::new(),
880 )?;
881 let mut updates: Vec<(String, String)> = Vec::new();
882 let mut links: BTreeMap<i64, BTreeSet<String>> = BTreeMap::new();
883 for i in 0..rs.len() {
884 let Some(Value::Str(sha)) = rs.get(i, "id") else {
885 continue;
886 };
887 let subject = match rs.get(i, "message") {
888 Some(Value::Str(s)) => s.as_str(),
889 _ => "",
890 };
891 let Some(number) = by_sha
892 .get(sha.as_str())
893 .copied()
894 .or_else(|| subject_pr(subject).filter(|n| known.contains(n)))
895 else {
896 continue;
897 };
898 let key = pr_key(number);
899 links.entry(number).or_default().insert(sha.clone());
900 if rs.get(i, "pr_id") != Some(&Value::Str(key.clone())) {
901 updates.push((sha.clone(), key));
902 }
903 }
904 updates.sort(); // by sha: the same store state writes the same records
905 for (sha, key) in updates {
906 w.set_prop(&sha, "pr_id", Value::Str(key))?;
907 }
908
909 let mut edges: Vec<(String, String, String)> = Vec::new();
910 for (number, shas) in links {
911 let src = pr_key(number);
912 let existing: BTreeSet<String> = w
913 .neighbors(&src, "MERGED_AS", Direction::Out)
914 .unwrap_or_default()
915 .into_iter()
916 .collect();
917 for sha in shas {
918 if !existing.contains(&sha) {
919 edges.push(("MERGED_AS".to_string(), src.clone(), sha));
920 }
921 }
922 }
923 if !edges.is_empty() {
924 w.ingest_with_edges("PR", Vec::new(), ingest, &edges)?;
925 }
926 Ok(())
927}
928
929/// Cypher behind [`file_state_from`] — the cumulative state of every live
930/// `File` node belonging to one unit.
931///
932/// The prefix filter is what keeps a submodule's files out of the parent's
933/// walk and vice versa. `startsWith` is the documented Cypher form (see
934/// `docs/site/query.md`); the parent unit's empty prefix matches everything, so
935/// its own submodules' keys are dropped in [`file_state_from`].
936const FILE_STATE_QUERY: &str =
937 "MATCH (f:File) WHERE startsWith(f.id, $prefix) RETURN f.id AS id, f.commits AS commits, \
938 f.n_commits AS n, f.top_author_id AS top, f.author_counts AS author_counts";
939
940/// Rebuild the in-memory per-file state from the `File` nodes already in the
941/// graph so incremental runs keep `commits` and ownership counts cumulative.
942///
943/// `nested` holds the key prefixes of any submodules inside this unit; their
944/// files belong to their own unit's walk and are dropped here.
945fn file_state_from(rs: &ResultSet, nested: &[String]) -> BTreeMap<String, FileState> {
946 let mut files = BTreeMap::new();
947 for i in 0..rs.len() {
948 let id = match rs.get(i, "id") {
949 Some(Value::Str(s)) => s.clone(),
950 _ => continue,
951 };
952 if nested.iter().any(|p| id.starts_with(p.as_str())) {
953 continue;
954 }
955 let mut st = FileState::default();
956 if let Some(Value::List(l)) = rs.get(i, "commits") {
957 st.commits = l
958 .iter()
959 .filter_map(|v| match v {
960 Value::Str(s) => Some(s.clone()),
961 _ => None,
962 })
963 .collect();
964 }
965 if let Some(Value::Int(n)) = rs.get(i, "n") {
966 st.n_commits = *n as usize;
967 }
968 match rs.get(i, "author_counts") {
969 Some(Value::List(l)) => st.set_author_counts(l),
970 // A node written before `author_counts` existed (a store built by
971 // 0.4.x). Fall back to the old approximation — the whole prior
972 // history credited to the incumbent — so those stores keep
973 // working. The prop is written on the next touch, from which point
974 // ownership tracks reality; a full re-ingest repairs it at once.
975 _ => {
976 if let Some(Value::Str(t)) = rs.get(i, "top") {
977 st.by_author.insert(t.clone(), st.n_commits.max(1));
978 }
979 }
980 }
981 files.insert(id, st);
982 }
983 files
984}
985
986/// The real per-name distribution for every author of this unit, read from the
987/// whole history rather than from the window.
988///
989/// A store built by 0.6.0 has `Author` nodes carrying a `name` and no
990/// `name_counts`, and nothing in the graph records which spelling each commit
991/// used — the display name it holds is the *first* one that version happened to
992/// see, which on a repository where one person commits under two names is the
993/// wrong one about as often as not. Guessing from it cannot recover: seeding the
994/// incumbent with the prior commit count makes the true majority wait for more
995/// new commits than the entire history it is already behind by.
996///
997/// So the log is walked in full, once, on the first run that meets such a node.
998/// The result covers every commit up to `head`, this run's window included, so
999/// the caller uses it *instead of* the window rather than on top of it.
1000/// A store this version wrote never reaches here.
1001fn legacy_name_counts(
1002 w: &WriteGuard<'_>,
1003 p: &Pending,
1004) -> Result<BTreeMap<String, AuthorState>, CliError> {
1005 let mut out: BTreeMap<String, AuthorState> = BTreeMap::new();
1006 let stale = p.log.iter().any(|c| {
1007 w.has_node(&c.author_email)
1008 && w.node_ref(&c.author_email)
1009 .and_then(|n| n.prop("name_counts"))
1010 .is_none()
1011 });
1012 let Some(head) = p.head.as_deref().filter(|_| stale) else {
1013 return Ok(out);
1014 };
1015 for c in read_log(&p.unit.path, None, head)? {
1016 out.entry(c.author_email).or_default().touch(&c.author_name);
1017 }
1018 Ok(out)
1019}
1020
1021/// Accumulated effect of one log window, before anything is written.
1022#[derive(Default)]
1023struct Walk {
1024 files: BTreeMap<String, FileState>,
1025 /// Email → the names its commits were authored under, and how many each.
1026 authors: BTreeMap<String, AuthorState>,
1027 commit_rows: Vec<BTreeMap<String, Value>>,
1028 touched_edges: Vec<(String, String, String)>,
1029 /// Paths whose `File` props changed in this window.
1030 dirty: BTreeSet<String>,
1031 deleted: BTreeSet<String>,
1032 /// Node renames to apply, collapsed across chained renames.
1033 renamed: Vec<(String, String)>,
1034 /// Any old path → its final path in this window, for edge retargeting.
1035 alias: BTreeMap<String, String>,
1036}
1037
1038impl Walk {
1039 fn rename(&mut self, from: &str, to: &str, node_exists: bool) {
1040 if let Some(e) = self.renamed.iter_mut().find(|(_, t)| t == from) {
1041 e.1 = to.to_string();
1042 } else if node_exists {
1043 // A node cannot be both deleted and moved; the move wins.
1044 self.deleted.remove(from);
1045 self.renamed.push((from.to_string(), to.to_string()));
1046 } else {
1047 self.deleted.insert(from.to_string());
1048 }
1049 for v in self.alias.values_mut() {
1050 if v == from {
1051 *v = to.to_string();
1052 }
1053 }
1054 self.alias.insert(from.to_string(), to.to_string());
1055 }
1056}
1057
1058/// The `GitSync` props that say *how* a unit was ingested, as opposed to how
1059/// far. Compared against the stored node so that changing a flag refreshes the
1060/// marker even when there is nothing new to walk.
1061fn marker_flag_props(unit: &RepoUnit, opts: &IngestGitOpts) -> Vec<(String, Value)> {
1062 vec![
1063 (
1064 "repo".into(),
1065 Value::Str(unit.path.to_string_lossy().into_owned()),
1066 ),
1067 ("recurse".into(), Value::Bool(opts.recurse_submodules)),
1068 ("prs".into(), Value::Bool(opts.prs)),
1069 ("structure".into(), Value::Bool(opts.structure)),
1070 ("docs".into(), Value::Bool(opts.docs)),
1071 ]
1072}
1073
1074/// One unit and everything read about it before the write pass opens.
1075struct Pending {
1076 unit: RepoUnit,
1077 /// The sha its marker resumes from, absent on a first run.
1078 since: Option<String>,
1079 /// The stored marker disagrees with this run's flags.
1080 stale: bool,
1081 /// `None` when the unit has no commits at all.
1082 head: Option<String>,
1083 log: Vec<GitCommit>,
1084 /// Key prefixes of the submodules nested inside this unit, whose files
1085 /// belong to their own walk.
1086 nested: Vec<String>,
1087}
1088
1089/// Whether `db_dir` would open faster if this run left a snapshot behind.
1090///
1091/// A full ingest always qualifies: it is the run that writes the whole history
1092/// into an empty WAL, and leaving that WAL to be replayed on every subsequent
1093/// open is the single largest fixed cost in the hook path. An incremental run
1094/// qualifies when the store has no snapshot at all — a store built by an older
1095/// release, or one whose full ingest predates this rule — or when the tail
1096/// past the last snapshot has grown beyond [`SNAPSHOT_WAL_BYTES`].
1097///
1098/// The tail *is* `wal.bin`: a snapshot replaces the live WAL with a minimal
1099/// baseline, so its length on disk measures exactly what an open has to replay.
1100///
1101/// Only a run that reached its write phase asks. A run with nothing to take
1102/// returns before entering a write scope, so a store with no snapshot and no
1103/// new commits keeps waiting for a run that has work — which in the hook path
1104/// is the next commit.
1105fn snapshot_due(db_dir: &Path, full: bool) -> bool {
1106 if full || !db_dir.join("snapshot.bin").exists() {
1107 return true;
1108 }
1109 std::fs::metadata(db_dir.join("wal.bin")).is_ok_and(|m| m.len() > SNAPSHOT_WAL_BYTES)
1110}
1111
1112/// Write a snapshot if one is due, after the run's own writes are committed.
1113///
1114/// The WAL is archived rather than dropped — see [`AUTOMATIC_SNAPSHOT`], which
1115/// every snapshot mushroomdb takes on its own shares — and every archive is
1116/// kept: [`AUTO_SNAPSHOT_RETENTION`] is `None` by default, so a store that
1117/// syncs on every commit grows an archive per 4 MiB of churn rather than
1118/// losing any of its history. `mushroomdb snapshot <db> --retention N` is how
1119/// a caller bounds that growth once they have decided the trade is worth it.
1120///
1121/// [`AUTOMATIC_SNAPSHOT`]: crate::AUTOMATIC_SNAPSHOT
1122/// [`AUTO_SNAPSHOT_RETENTION`]: crate::AUTO_SNAPSHOT_RETENTION
1123///
1124/// # Why it never fails the run
1125///
1126/// A snapshot is an optimisation for the *next* open, and the data is already
1127/// durable either way:
1128///
1129/// * `Busy` — another process holds the cross-process write lock. Skipped
1130/// without a word; the next run that qualifies takes it.
1131/// * anything else — reported on stderr, because a store that cannot be
1132/// snapshotted is worth knowing about, and the run still succeeds.
1133///
1134/// Call this only from a run that reached its write phase. A run with nothing
1135/// to do returns before entering a write scope, and taking the lock to
1136/// snapshot would turn a genuine no-op into a write.
1137fn snapshot_if_due(db: &SharedDb, db_dir: &Path, full: bool) {
1138 if !snapshot_due(db_dir, full) {
1139 return;
1140 }
1141 let taken = db
1142 .write_with_wait(WRITE_LOCK_WAIT)
1143 .and_then(|mut g| crate::snapshot_automatically(&mut g));
1144 match taken {
1145 Ok(()) | Err(GraphError::Busy { .. }) => {}
1146 Err(e) => eprintln!("snapshot skipped: {e}"),
1147 }
1148}
1149
1150pub fn run_ingest_git(db_dir: &Path, opts: &IngestGitOpts) -> Result<IngestGitReport, CliError> {
1151 // Absolute and symlink-free: it is recorded on the marker, and a later run
1152 // has no reason to share this one's working directory.
1153 let repo = canonical(&opts.repo);
1154 let units = repo_units(&repo, opts.recurse_submodules)?;
1155 let prefixes: Vec<String> = units
1156 .iter()
1157 .map(|u| u.prefix.clone())
1158 .filter(|p| !p.is_empty())
1159 .collect();
1160
1161 // The same gate `serve` and `mcp` use, read before `open` creates the
1162 // directory's files: after the open it would always answer "existing".
1163 // A run that finds nothing to write leaves the WAL empty, which this
1164 // still reads as new, so the schema lands with the first run that writes.
1165 let is_new_store = !core_api::restore::holds_a_store(db_dir);
1166 let db = SharedDb::open(db_dir)?;
1167 let mut report = IngestGitReport {
1168 submodules: units.len() - 1,
1169 ..Default::default()
1170 };
1171 if opts.ensure_gitignore {
1172 report.gitignore_added = ensure_gitignore(&repo, db_dir)?;
1173 }
1174
1175 let mut pending: Vec<Pending> = Vec::new();
1176 {
1177 let r = db.read();
1178 for unit in units {
1179 let nested = prefixes
1180 .iter()
1181 .filter(|p| **p != unit.prefix && p.starts_with(&unit.prefix))
1182 .cloned()
1183 .collect();
1184 let (since, stale) = match r.node_ref(&unit.sync_key) {
1185 Some(n) => (
1186 match n.prop("sha") {
1187 Some(Value::Str(s)) if !s.is_empty() => Some(s),
1188 _ => None,
1189 },
1190 marker_flag_props(&unit, opts)
1191 .into_iter()
1192 .any(|(k, v)| n.prop(&k) != Some(v)),
1193 ),
1194 None => (None, false),
1195 };
1196 pending.push(Pending {
1197 unit,
1198 since,
1199 stale,
1200 head: None,
1201 log: Vec::new(),
1202 nested,
1203 });
1204 }
1205 }
1206 // The repository itself is always the first unit, and it is what "this
1207 // database has seen this repo before" means.
1208 report.incremental = pending[0].since.is_some();
1209
1210 for p in &mut pending {
1211 // Pin the end of the walk and the marker to one sha, resolved first. A
1212 // commit landing mid-run then falls outside this range instead of being
1213 // skipped by a marker that advanced past it.
1214 let Some(head) = head_sha(&p.unit.path)? else {
1215 continue; // this unit has no commits yet
1216 };
1217 let mut log = read_log(&p.unit.path, p.since.as_deref(), &head)?;
1218 localise(&mut log, &p.unit.prefix, &gitlink_paths(&p.unit.path));
1219 p.head = Some(head);
1220 p.log = log;
1221 }
1222
1223 let prs = if opts.prs {
1224 fetch_prs(&repo)
1225 } else {
1226 Vec::new()
1227 };
1228 if prs.is_empty() && pending.iter().all(|p| p.log.is_empty() && !p.stale) {
1229 // Nothing new: leave the store untouched so `commit_seq` does not move.
1230 return Ok(report);
1231 }
1232
1233 // Held for the rest of the run. Another process holding the store's
1234 // cross-process write lock is a retry, not a failure: nothing was written.
1235 let mut w = db.write_with_wait(WRITE_LOCK_WAIT).map_err(|e| match e {
1236 GraphError::Busy { .. } => CliError(BUSY_MESSAGE.to_string()),
1237 other => CliError(other.to_string()),
1238 })?;
1239 let ingest = IngestOptions::default(); // key `id`, auto-FK suffix `_id`
1240
1241 // A store this run is creating is declared on before its first write, so a
1242 // failure part-way through still leaves it searchable. An existing store
1243 // is never re-declared on: declaring full-text rebuilds its index at every
1244 // later open (ledger row 36), and only `schema apply` may ask for that.
1245 if is_new_store {
1246 let diff = w.apply_schema(&structure::ingest_git_schema())?;
1247 report.rules_created.extend(
1248 diff.created
1249 .iter()
1250 .filter_map(|entry| entry.strip_prefix("rule:"))
1251 .map(str::to_string),
1252 );
1253 }
1254
1255 // Pull requests first: their nodes are what `Commit.pr_id` resolves to.
1256 if !prs.is_empty() {
1257 ingest_prs(&mut w, &prs, &ingest, &mut report)?;
1258 }
1259
1260 let mut authors: BTreeSet<String> = BTreeSet::new();
1261 let mut work = StructureWork::default();
1262 for p in &pending {
1263 if !p.log.is_empty() {
1264 ingest_unit(
1265 &mut w,
1266 p,
1267 opts,
1268 &ingest,
1269 &mut report,
1270 &mut authors,
1271 &mut work,
1272 )?;
1273 }
1274 }
1275 report.authors = authors.len();
1276
1277 // Rules and fulltext, created after the data so each backfills once. They
1278 // span every unit, so they are declared once for the database rather than
1279 // once per repository walked.
1280 if report.commits > 0 && !w.rules().iter().any(|r| r.name == "co_changed") {
1281 let co = Predicate::Overlap {
1282 field: "commits".into(),
1283 min: CO_CHANGE_MIN,
1284 };
1285 w.create_rule(RuleDef {
1286 name: "co_changed".into(),
1287 src_label: "File".into(),
1288 dst_label: "File".into(),
1289 predicate: co.clone(),
1290 edge_type: "CO_CHANGED".into(),
1291 weight_prop: Some("score".into()),
1292 max_edges: Some(10),
1293 approximate: false,
1294 via_label: None,
1295 via_edge: None,
1296 via_dir: None,
1297 namespace: None,
1298 })?;
1299 w.create_rule(RuleDef {
1300 name: "knows".into(),
1301 src_label: "Author".into(),
1302 dst_label: "File".into(),
1303 predicate: co,
1304 edge_type: "KNOWS".into(),
1305 weight_prop: Some("score".into()),
1306 max_edges: Some(20),
1307 approximate: false,
1308 via_label: Some("File".into()),
1309 via_edge: Some("TOP_AUTHOR".into()),
1310 via_dir: Some(Direction::In),
1311 namespace: None,
1312 })?;
1313 report
1314 .rules_created
1315 .extend(["co_changed".to_string(), "knows".to_string()]);
1316 for (l, field) in [("File", "path"), ("Commit", "message"), ("Author", "name")] {
1317 if !w
1318 .fulltext_pairs()
1319 .contains(&(l.to_string(), field.to_string()))
1320 {
1321 w.enable_fulltext(l, field)?;
1322 }
1323 }
1324 }
1325
1326 // The working tree, on top of the history. It runs after the commit walk
1327 // so every `File` node it reads exists. On a store this run created, its
1328 // rules came with the schema declared before the first write.
1329 if opts.structure {
1330 // A first run has nothing to be incremental against, and a run whose
1331 // flags changed (structure or docs just turned on) has to revisit
1332 // files its predecessor deliberately skipped.
1333 let full = !report.incremental || pending.iter().any(|p| p.stale);
1334 report.structure = if full {
1335 structure::refresh_all(&mut w, &repo, "", opts.docs)?
1336 } else {
1337 let mut paths = work.touched;
1338 paths.extend(structure::importers_of(&w, &work.stale)?);
1339 // The retired keys themselves, not just the files that named them.
1340 // An incremental refresh sweeps orphaned symbols only under the
1341 // paths it is handed, and a file this window deleted or renamed
1342 // away is exactly where the orphans are.
1343 paths.extend(work.stale.iter().cloned());
1344 let paths: Vec<String> = paths.into_iter().collect();
1345 structure::refresh_files(&mut w, &repo, "", &paths, opts.docs)?
1346 };
1347 // A store this run did not create — `mcp` made it, or an older run
1348 // with `--no-structure` did — gets ingest-git's own rules and text
1349 // fields, declared after the props so each rule backfills once. Only
1350 // what is absent, so a re-run declares nothing; never the memory
1351 // defaults, which reach an existing store only through `schema apply`.
1352 if !is_new_store {
1353 let missing = structure::missing_structure_schema(&w);
1354 let diff = w.apply_schema(&missing)?;
1355 report.rules_created.extend(
1356 diff.created
1357 .iter()
1358 .filter_map(|entry| entry.strip_prefix("rule:"))
1359 .map(str::to_string),
1360 );
1361 }
1362 }
1363
1364 // Every commit this run could link is in the graph by now.
1365 if !prs.is_empty() {
1366 link_prs(&mut w, &prs, &ingest)?;
1367 }
1368
1369 // The markers go last, once every phase of the run has succeeded. They say
1370 // how far a *complete* run got, so a failure anywhere above — a working-tree
1371 // batch that will not commit, a `gh` link pass that errors — leaves them
1372 // where they were and the next run re-walks the same window rather than
1373 // stepping over it. Re-walking a window that was partly applied is safe:
1374 // commits already in the graph are skipped as duplicate keys, file props
1375 // are rewritten from the recomputed state, and a rename whose node already
1376 // moved finds nothing to move.
1377 for p in &pending {
1378 write_marker(&mut w, p, opts)?;
1379 }
1380
1381 // Every write this run makes is now committed, so leave the store in the
1382 // shape that opens fastest. The guard is released first: `snapshot()` takes
1383 // the same cross-process write lock this one holds.
1384 drop(w);
1385 snapshot_if_due(&db, db_dir, !report.incremental);
1386 Ok(report)
1387}
1388
1389/// Wall-clock seconds since the Unix epoch, for [`SYNCED_AT`].
1390///
1391/// This is the one place the ingest reads a clock. A clock before the epoch
1392/// reads as `0` rather than going negative.
1393fn now_unix() -> i64 {
1394 std::time::SystemTime::now()
1395 .duration_since(std::time::UNIX_EPOCH)
1396 .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
1397}
1398
1399/// Marker prop holding when this store last took data from the repository.
1400///
1401/// `Commit.ts` says when the work was *written*, which on a store synced to
1402/// its repository's head is indistinguishable from now — so it cannot answer
1403/// "how stale is my graph". This can. Absent on a store built before it
1404/// existed, which readers must tolerate.
1405pub const SYNCED_AT: &str = "synced_at";
1406
1407/// Record how far this unit got, and under what flags.
1408///
1409/// Props are written only where they differ, so a run that touches one unit
1410/// does not churn the markers of the others — and a run that changes nothing
1411/// writes nothing at all, [`SYNCED_AT`] included.
1412///
1413/// The stamp therefore means "when this run last changed this marker": a new
1414/// head, or a flag ingested differently from last time. A no-op re-run leaving
1415/// it alone is the point rather than a gap — a graph that had nothing to take
1416/// is not staler for having checked.
1417fn write_marker(w: &mut WriteGuard<'_>, p: &Pending, opts: &IngestGitOpts) -> Result<(), CliError> {
1418 let Some(head) = p.head.as_deref() else {
1419 return Ok(()); // no commits, so nothing to resume from
1420 };
1421 let mut props = marker_flag_props(&p.unit, opts);
1422 if !p.log.is_empty() {
1423 props.push(("sha".into(), Value::Str(head.to_string())));
1424 }
1425 let key = p.unit.sync_key.clone();
1426 if w.has_node(&key) {
1427 let mut changed = false;
1428 for (k, v) in props {
1429 let current = w.node_ref(&key).and_then(|n| n.prop(&k));
1430 if current.as_ref() != Some(&v) {
1431 w.set_prop(&key, &k, v)?;
1432 changed = true;
1433 }
1434 }
1435 if changed {
1436 w.set_prop(&key, SYNCED_AT, Value::Int(now_unix()))?;
1437 }
1438 return Ok(());
1439 }
1440 if p.log.is_empty() {
1441 return Ok(());
1442 }
1443 // It carries `id` like every other label here, so the key is readable from
1444 // Cypher.
1445 props.push(("id".into(), Value::Str(key.clone())));
1446 props.push((SYNCED_AT.into(), Value::Int(now_unix())));
1447 props.sort_by(|a, b| a.0.cmp(&b.0));
1448 w.insert_node("GitSync", &key, props)?;
1449 Ok(())
1450}
1451
1452/// Walk one unit's new commits into the graph.
1453///
1454/// Every path in `p.log` is already a repository-wide key, so this is the
1455/// single-repository algorithm unchanged: only the `File` state it starts from
1456/// and the counts it adds to are scoped to the unit.
1457fn ingest_unit(
1458 w: &mut WriteGuard<'_>,
1459 p: &Pending,
1460 opts: &IngestGitOpts,
1461 ingest: &IngestOptions,
1462 report: &mut IngestGitReport,
1463 authors: &mut BTreeSet<String>,
1464 work: &mut StructureWork,
1465) -> Result<(), CliError> {
1466 let log = &p.log;
1467 let incremental = p.since.is_some();
1468
1469 let mut walk = Walk {
1470 files: if incremental {
1471 let params =
1472 BTreeMap::from([("prefix".to_string(), Value::Str(p.unit.prefix.clone()))]);
1473 file_state_from(&w.query(FILE_STATE_QUERY, ¶ms)?, &p.nested)
1474 } else {
1475 BTreeMap::new()
1476 },
1477 ..Default::default()
1478 };
1479
1480 for c in log {
1481 walk.authors
1482 .entry(c.author_email.clone())
1483 .or_default()
1484 .touch(&c.author_name);
1485 walk.commit_rows.push(BTreeMap::from([
1486 ("id".to_string(), Value::Str(c.sha.clone())),
1487 ("message".to_string(), Value::Str(c.subject.clone())),
1488 ("ts".to_string(), Value::Int(c.ts)),
1489 ("author_id".to_string(), Value::Str(c.author_email.clone())),
1490 ]));
1491 for ch in &c.changes {
1492 match ch {
1493 Change::Added(p) | Change::Modified(p) => {
1494 if path_excluded(p, &opts.exclude) {
1495 continue;
1496 }
1497 walk.deleted.remove(p);
1498 walk.files.entry(p.clone()).or_default().touch(
1499 &c.sha,
1500 &c.author_email,
1501 opts.max_commits_per_file,
1502 );
1503 walk.dirty.insert(p.clone());
1504 walk.touched_edges
1505 .push(("TOUCHED".into(), c.sha.clone(), p.clone()));
1506 }
1507 Change::Deleted(p) => {
1508 if path_excluded(p, &opts.exclude) {
1509 continue;
1510 }
1511 walk.files.remove(p);
1512 walk.dirty.remove(p);
1513 walk.deleted.insert(p.clone());
1514 }
1515 Change::Renamed { from, to } => {
1516 if path_excluded(to, &opts.exclude) {
1517 // Moved out of scope: drop the old node, keep no alias
1518 // so its TOUCHED edges are filtered out below.
1519 walk.files.remove(from);
1520 walk.dirty.remove(from);
1521 walk.deleted.insert(from.clone());
1522 continue;
1523 }
1524 let mut st = walk.files.remove(from).unwrap_or_default();
1525 st.touch(&c.sha, &c.author_email, opts.max_commits_per_file);
1526 walk.files.insert(to.clone(), st);
1527 walk.dirty.remove(from);
1528 walk.dirty.insert(to.clone());
1529 walk.deleted.remove(to);
1530 walk.touched_edges
1531 .push(("TOUCHED".into(), c.sha.clone(), to.clone()));
1532 let exists = incremental && w.has_node(from);
1533 walk.rename(from, to, exists);
1534 }
1535 }
1536 }
1537 }
1538
1539 // What the working-tree pass has to look at afterwards: the paths this
1540 // window changed, and the keys it left pointing at nothing. A key that
1541 // ended the window live again — a file moved away and back — is neither.
1542 work.touched.extend(walk.dirty.iter().cloned());
1543 for key in walk.deleted.iter().chain(walk.alias.keys()) {
1544 if !walk.files.contains_key(key) {
1545 work.stale.insert(key.clone());
1546 }
1547 }
1548
1549 // 1. Authors first: the auto-FK rules for `Commit.author_id` and
1550 // `File.top_author_id` only infer once their targets resolve to Author.
1551 // An email already in the graph keeps its accumulated `name_counts`, so
1552 // the display name is the majority spelling over the whole history
1553 // rather than over this window.
1554 let legacy = legacy_name_counts(w, p)?;
1555 let mut author_rows = Vec::new();
1556 let mut author_updates = Vec::new();
1557 for (email, seen) in &walk.authors {
1558 let mut merged = AuthorState::default();
1559 match (
1560 w.node_ref(email).and_then(|n| n.prop("name_counts")),
1561 legacy.get(email),
1562 ) {
1563 // The usual path: what the graph already counted, plus this window.
1564 (Some(Value::List(l)), _) => {
1565 merged.set_name_counts(&l);
1566 merged.absorb(seen);
1567 }
1568 // A node written before `name_counts` existed. The recovered walk
1569 // is the whole history and already contains this window, so the
1570 // window is not added to it.
1571 (_, Some(whole_history)) => merged = whole_history.clone(),
1572 // A new author, or a unit whose history was already recovered.
1573 _ => merged.absorb(seen),
1574 }
1575 let props = [
1576 ("name", Value::Str(merged.display_name())),
1577 ("name_counts", merged.name_counts_value()),
1578 ];
1579 if w.has_node(email) {
1580 for (field, want) in props {
1581 if w.node_ref(email).and_then(|n| n.prop(field)).as_ref() != Some(&want) {
1582 author_updates.push((email.clone(), field, want));
1583 }
1584 }
1585 } else {
1586 let mut row = BTreeMap::from([("id".to_string(), Value::Str(email.clone()))]);
1587 row.extend(props.into_iter().map(|(f, v)| (f.to_string(), v)));
1588 author_rows.push(row);
1589 }
1590 }
1591 if !author_updates.is_empty() {
1592 let mut b = w.batch();
1593 for (key, field, value) in author_updates {
1594 b.set_prop(&key, field, value);
1595 }
1596 b.commit()?;
1597 }
1598 let a = w.ingest_with_edges("Author", author_rows, ingest, &[])?;
1599 report.rules_created.extend(a.rules_created);
1600 authors.extend(walk.authors.keys().cloned());
1601
1602 // 2. Deletes run first so a rename can claim a path freed in this same
1603 // window, then renames carry each node (and its history) to its new path.
1604 // `walk.files` is the authority on what is still live: a rename whose
1605 // destination is not in it is a delete, not a move.
1606 for p in &walk.deleted {
1607 if w.has_node(p) {
1608 w.delete_node(p)?;
1609 report.deleted += 1;
1610 }
1611 }
1612 for (from, to) in &walk.renamed {
1613 if !w.has_node(from) || from == to {
1614 // Nothing to move, or a rename that swapped back to its own path.
1615 continue;
1616 }
1617 if !walk.files.contains_key(to) {
1618 // The destination did not survive the window — it was deleted, or
1619 // moved into an excluded path, after this rename. The node goes
1620 // with it; renaming into a dead path would strand a phantom node
1621 // that no later phase refreshes.
1622 //
1623 // This is an eviction, not a delete: no path this window deleted
1624 // ever named this node. Counting it as a delete overstated how much
1625 // the window removed, which is what §5.12 splits.
1626 w.delete_node(from)?;
1627 report.evicted += 1;
1628 continue;
1629 }
1630 if w.has_node(to) {
1631 // A pre-existing node already holds the destination path (deleted
1632 // earlier in this window, then claimed by this rename).
1633 //
1634 // Counted as a delete, deliberately: unlike the eviction above,
1635 // this node's own path *was* removed in this window — the rename
1636 // only decides who claims it next.
1637 w.delete_node(to)?;
1638 report.deleted += 1;
1639 }
1640 w.rename_node(from, to)?;
1641 // The key moved, so the `id` prop must move with it. Phase 3 also sets
1642 // it for every dirty path; doing it here keeps the invariant local to
1643 // the rename and independent of that filter.
1644 w.set_prop(to, "id", Value::Str(to.clone()))?;
1645 report.renamed += 1;
1646 }
1647
1648 // 3. File nodes. Existing nodes are updated in place (including `id`, which
1649 // must follow the key after a rename); new paths go through ingest so
1650 // the `top_author_id` auto-FK rule is inferred.
1651 // Updates to existing nodes go in one batch rather than one WAL commit
1652 // per property. Every frame in the WAL is replayed on every later open,
1653 // and each one re-fires the rules watching the property it carries, so a
1654 // run that appends a frame per property makes every subsequent open
1655 // slower for as long as that frame lives. The op order inside the batch
1656 // is the order the individual writes had, so the resulting state is the
1657 // same one either form produces.
1658 let mut new_file_rows = Vec::new();
1659 let mut updates = Vec::new();
1660 let mut written = 0usize;
1661 for (path, st) in &walk.files {
1662 if incremental && !walk.dirty.contains(path) {
1663 continue;
1664 }
1665 written += 1;
1666 let props = file_props(st, path);
1667 if w.has_node(path) {
1668 updates.extend(props.into_iter().map(|(k, v)| (path.clone(), k, v)));
1669 } else {
1670 new_file_rows.push(props.into_iter().collect::<BTreeMap<_, _>>());
1671 }
1672 }
1673 if !updates.is_empty() {
1674 let mut b = w.batch();
1675 for (key, field, value) in updates {
1676 b.set_prop(&key, &field, value);
1677 }
1678 b.commit()?;
1679 }
1680 let f = w.ingest_with_edges("File", new_file_rows, ingest, &[])?;
1681 report.rules_created.extend(f.rules_created);
1682 // What this run wrote, not every file it has ever seen: an incremental run
1683 // loads the whole known file set to fold the new commits into it, and
1684 // reporting that set made a one-file run print the size of the repository.
1685 report.files += written;
1686
1687 // 4. Commits, then their TOUCHED edges. The two must be separate batches:
1688 // a batch that both inserts nodes firing a new rule and carries a user
1689 // edge of a not-yet-interned type writes a WAL frame that cannot be
1690 // replayed (`Intern` records are emitted in a pre-pass, but on replay the
1691 // rule fires — and interns its edge type — before the later `Intern`
1692 // record is read). See the report for a reproducer.
1693 let c = w.ingest_with_edges("Commit", walk.commit_rows, ingest, &[])?;
1694 report.rules_created.extend(c.rules_created);
1695 report.commits += log.len();
1696
1697 // Edges name File keys, so files must already exist. A path renamed later
1698 // in this same window is retargeted to where its node ended up.
1699 let touched: Vec<(String, String, String)> = walk
1700 .touched_edges
1701 .into_iter()
1702 .map(|(t, sha, p)| {
1703 let p = walk.alias.get(&p).cloned().unwrap_or(p);
1704 (t, sha, p)
1705 })
1706 .filter(|(_, _, p)| walk.files.contains_key(p))
1707 .collect();
1708 if !touched.is_empty() {
1709 w.ingest_with_edges("Commit", Vec::new(), ingest, &touched)?;
1710 }
1711 Ok(())
1712}
1713
1714pub fn format_ingest_git(r: &IngestGitReport) -> String {
1715 let mut out = format!(
1716 "ingest-git: {} commit(s), {} file(s), {} author(s){}\n",
1717 r.commits,
1718 r.files,
1719 r.authors,
1720 if r.incremental { " (incremental)" } else { "" }
1721 );
1722 if r.renamed + r.deleted + r.evicted > 0 {
1723 out.push_str(&format!(
1724 " renamed {} deleted {} evicted {}\n",
1725 r.renamed, r.deleted, r.evicted
1726 ));
1727 }
1728 if r.submodules + r.prs > 0 {
1729 out.push_str(&format!(
1730 " submodules {} pull requests {}\n",
1731 r.submodules, r.prs
1732 ));
1733 }
1734 let s = &r.structure;
1735 if s.files_scanned > 0 {
1736 out.push_str(&format!(
1737 " scanned {} file(s): {} symbol(s), {} import(s), {} call(s), {} mention(s)\n",
1738 s.files_scanned, s.symbols, s.imports, s.calls, s.mentions
1739 ));
1740 if s.skipped_large + s.symbols_capped > 0 {
1741 out.push_str(&format!(
1742 " hash-only {} symbol cap hit on {}\n",
1743 s.skipped_large, s.symbols_capped
1744 ));
1745 }
1746 }
1747 if r.gitignore_added {
1748 out.push_str(" added the database directory to .gitignore\n");
1749 }
1750 if !r.rules_created.is_empty() {
1751 out.push_str(&format!(" rules: {}\n", r.rules_created.join(", ")));
1752 }
1753 out
1754}
1755
1756#[cfg(test)]
1757mod tests {
1758 use super::*;
1759
1760 #[test]
1761 fn exclude_matches_prefix_extension_and_substring() {
1762 let pats = vec![
1763 "target/".to_string(),
1764 "*.lock".into(),
1765 "node_modules".into(),
1766 ];
1767 assert!(path_excluded("target/debug/foo.rs", &pats));
1768 assert!(
1769 !path_excluded("targeted/foo.rs", &pats),
1770 "prefix needs the slash"
1771 );
1772 assert!(path_excluded("Cargo.lock", &pats));
1773 assert!(!path_excluded("Cargo.toml", &pats));
1774 assert!(path_excluded("ui/node_modules/x/y.js", &pats));
1775 assert!(!path_excluded("src/lib.rs", &pats));
1776 assert!(!path_excluded("anything", &[]));
1777 }
1778
1779 /// A `*.` pattern is a file-name suffix, so a compound one works. Reading
1780 /// only the last dot segment would make `*.min.js` — a default — inert,
1781 /// and generated bundles are both the largest files in a tree and the ones
1782 /// least worth parsing.
1783 #[test]
1784 fn a_compound_suffix_pattern_matches() {
1785 let defaults: Vec<String> = DEFAULT_EXCLUDES.iter().map(|p| (*p).to_string()).collect();
1786 // Not under `dist/`, so only the suffix rule can match it.
1787 assert!(path_excluded("ui/build/bundle.min.js", &defaults));
1788 assert!(path_excluded("bundle.min.js", &defaults));
1789 assert!(
1790 !path_excluded("ui/src/app.js", &defaults),
1791 "an ordinary source file is not a bundle"
1792 );
1793 assert!(!path_excluded("ui/src/minify.js", &defaults));
1794 // The single-extension form is unchanged, and a bare suffix is not a
1795 // match: `*.lock` means something *dot* lock.
1796 assert!(path_excluded("Cargo.lock", &defaults));
1797 assert!(!path_excluded(".lock", &defaults));
1798 assert!(!path_excluded("src/lib.rs", &defaults));
1799 }
1800
1801 #[test]
1802 fn file_props_split_dir_and_ext() {
1803 let mut st = FileState::default();
1804 st.touch("sha1", "a@x.test", 200);
1805 let m: BTreeMap<_, _> = file_props(&st, "src/a/b.rs").into_iter().collect();
1806 assert_eq!(m["dir"], Value::Str("src/a".into()));
1807 assert_eq!(m["ext"], Value::Str("rs".into()));
1808 assert_eq!(m["n_commits"], Value::Int(1));
1809 assert_eq!(m["id"], Value::Str("src/a/b.rs".into()));
1810 assert_eq!(m["top_author_id"], Value::Str("a@x.test".into()));
1811 assert_eq!(
1812 m["author_counts"],
1813 Value::List(vec![Value::Str("a@x.test\t1".into())])
1814 );
1815 let m: BTreeMap<_, _> = file_props(&FileState::default(), "README")
1816 .into_iter()
1817 .collect();
1818 assert_eq!(m["dir"], Value::Str(String::new()));
1819 assert_eq!(m["ext"], Value::Str(String::new()));
1820 }
1821
1822 /// The prop is the whole point of the incremental fix: it must survive a
1823 /// round trip so the next run resumes the real distribution, not the
1824 /// incumbent's total.
1825 #[test]
1826 fn author_counts_round_trip_preserves_the_distribution() {
1827 let mut st = FileState::default();
1828 for _ in 0..3 {
1829 st.touch("s", "alice@x.test", 200);
1830 }
1831 for _ in 0..4 {
1832 st.touch("s", "bob@x.test", 200);
1833 }
1834 let Value::List(encoded) = st.author_counts_value() else {
1835 panic!("author_counts must be a list");
1836 };
1837 assert_eq!(
1838 encoded,
1839 vec![
1840 Value::Str("alice@x.test\t3".into()),
1841 Value::Str("bob@x.test\t4".into()),
1842 ],
1843 "email order, so the prop is stable across runs"
1844 );
1845 let mut reloaded = FileState::default();
1846 reloaded.set_author_counts(&encoded);
1847 assert_eq!(reloaded.by_author, st.by_author);
1848 assert_eq!(reloaded.top_author(), "bob@x.test");
1849 }
1850
1851 /// Malformed entries are skipped, not fatal: the run degrades to the counts
1852 /// it can read rather than refusing to sync.
1853 #[test]
1854 fn author_counts_skips_entries_it_cannot_parse() {
1855 let mut st = FileState::default();
1856 st.set_author_counts(&[
1857 Value::Str("alice@x.test\t2".into()),
1858 Value::Str("no-tab-here".into()),
1859 Value::Str("bob@x.test\tnotanumber".into()),
1860 Value::Str("\t5".into()),
1861 Value::Int(7),
1862 ]);
1863 assert_eq!(st.by_author, BTreeMap::from([("alice@x.test".into(), 2)]));
1864 }
1865
1866 #[test]
1867 fn commits_list_is_capped_and_top_author_is_deterministic() {
1868 let mut st = FileState::default();
1869 for i in 0..5 {
1870 st.touch(&format!("sha{i}"), "b@x.test", 3);
1871 }
1872 st.touch("shaX", "a@x.test", 3);
1873 assert_eq!(st.commits, vec!["sha3", "sha4", "shaX"]);
1874 assert_eq!(st.n_commits, 6);
1875 assert_eq!(st.top_author(), "b@x.test");
1876
1877 // Past the cap the stored `n_commits` is the true total, not the length
1878 // of the truncated `commits` list, and `author_counts` sums to the same
1879 // number. A file over `--max-commits-per-file` would otherwise report a
1880 // history frozen at the cap.
1881 let m: BTreeMap<_, _> = file_props(&st, "src/hot.rs").into_iter().collect();
1882 assert_eq!(m["n_commits"], Value::Int(6));
1883 let Value::List(commits) = &m["commits"] else {
1884 panic!("commits must be a list");
1885 };
1886 assert_eq!(commits.len(), 3, "the list is still capped at 3");
1887 assert!(
1888 matches!(m["n_commits"], Value::Int(n) if n as usize > commits.len()),
1889 "n_commits must exceed the capped list once the cap is passed"
1890 );
1891 assert_eq!(
1892 m["author_counts"],
1893 Value::List(vec![
1894 Value::Str("a@x.test\t1".into()),
1895 Value::Str("b@x.test\t5".into()),
1896 ]),
1897 "the per-author counts sum to n_commits, not to the capped list"
1898 );
1899
1900 let mut tie = FileState::default();
1901 tie.touch("s", "b@x.test", 10);
1902 tie.touch("s", "a@x.test", 10);
1903 assert_eq!(tie.top_author(), "a@x.test", "ties break on smallest email");
1904 }
1905}