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