Skip to main content

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, &params)?, &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}