pub mod at_rev;
mod clones_head;
mod complexity_head;
pub(crate) mod consumer;
mod grouping;
mod imports_head;
mod lineage;
pub use at_rev::{ingest_complexity_at_rev, materialize_imports_at_rev};
pub use grouping::{apply_grouping, materialize_changes_bucketed};
pub use lineage::{materialize_changes_lineage, materialize_path_lineage};
use crossbeam_channel::bounded;
use super::FactsDb;
use crate::repo::Repo;
use crate::{CodeLoreError, CommitEvent, Options, Result};
use consumer::ingest_loop;
const DEFAULT_CHANNEL_CAPACITY: usize = 64;
static CHANNEL_CAPACITY_OVERRIDE: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub fn set_channel_capacity_override(n: usize) {
CHANNEL_CAPACITY_OVERRIDE.store(n, std::sync::atomic::Ordering::SeqCst);
}
fn channel_capacity() -> usize {
let n = CHANNEL_CAPACITY_OVERRIDE.load(std::sync::atomic::Ordering::SeqCst);
if n == 0 { DEFAULT_CHANNEL_CAPACITY } else { n }
}
#[derive(Debug, Default)]
pub struct IngestStats {
pub commits_ingested: usize,
pub changes_ingested: usize,
pub clones_ingested: usize,
}
impl FactsDb {
pub fn ingest<R: Repo>(&self, repo: &R, opts: &Options) -> Result<IngestStats> {
if opts.head_only_ingest {
return self.ingest_head_only(repo, opts);
}
let team_map_path = opts
.team_map_file
.clone()
.or_else(|| crate::identity::discover_team_map(&opts.repo_path));
let team_map = crate::identity::team_map::load(team_map_path.as_deref())?;
let bot_patterns = crate::identity::BotPatterns::from_repo(&opts.repo_path);
let (tx, rx) = bounded::<CommitEvent>(channel_capacity());
let stats = std::thread::scope(|s| -> Result<IngestStats> {
let producer = s.spawn(|| -> Result<()> {
let walk = repo.walk_commits(opts)?;
for event in walk {
let event = event?;
tx.send(event)
.map_err(|e| CodeLoreError::Analysis(format!("channel send: {e}")))?;
}
drop(tx); Ok(())
});
let paths_filter = crate::paths_filter::PathsFilter::from_opts(opts)?;
let stats = ingest_loop(self, rx, &team_map, &bot_patterns, &paths_filter)?;
let join_result = producer
.join()
.map_err(|payload| CodeLoreError::Repo(format_panic_payload(&payload)))?;
join_result?;
Ok(stats)
})?;
match future_dated_commit_warning(self) {
Ok(Some(msg)) => tracing::warn!("{msg}"),
Ok(None) => {}
Err(e) => tracing::debug!("future-date ingest check skipped: {e}"),
}
let head_rev = current_head_rev(self)?;
let live_paths = query_live_paths(self)?;
self.ingest_complexity_at_head(repo, opts, &live_paths, &head_rev)?;
crate::kamei::enrich(self, opts.use_canonical_lineage)?;
let clones_n = self.populate_clones_at_head(repo, opts, &live_paths, &head_rev)?;
let imports_n = self.populate_imports_at_head(repo, opts, &live_paths, &head_rev)?;
tracing::info!("imports: {imports_n} edges ingested at HEAD across Tier-1 source files");
let resolved_n = self.resolve_imports_at_head(&live_paths, &head_rev)?;
tracing::info!(
"imports: {resolved_n} of {imports_n} import edges resolved to tracked paths"
);
if let Some(group_file) = opts.group_file.as_ref() {
let group_map = super::groups::GroupMap::from_file(group_file, opts.strict_grouping)
.map_err(|e| match e {
super::groups::GroupParseError::Io(io) => {
CodeLoreError::RepoIo(std::io::Error::new(
io.kind(),
format!("read group-file {}: {io}", group_file.display()),
))
}
other => CodeLoreError::Analysis(format!("--group-file: {other}")),
})?;
apply_grouping(self, &group_map)?;
}
let mut stats = stats;
stats.clones_ingested = clones_n;
Ok(stats)
}
fn ingest_head_only<R: Repo>(&self, repo: &R, opts: &Options) -> Result<IngestStats> {
let head_rev = repo.head_sha()?;
let paths_filter = crate::paths_filter::PathsFilter::from_opts(opts)?;
let live_paths: Vec<String> = repo
.tracked_paths_at_head()?
.into_iter()
.filter(|p| {
let rel_path = std::path::Path::new(p);
!crate::paths_filter::is_git_metadata(rel_path)
&& !paths_filter.is_excluded(rel_path, false)
})
.collect();
self.ingest_complexity_at_head(repo, opts, &live_paths, &head_rev)?;
let imports_n = self.populate_imports_at_head(repo, opts, &live_paths, &head_rev)?;
let resolved_n = self.resolve_imports_at_head(&live_paths, &head_rev)?;
tracing::info!(
"imports: {resolved_n} of {imports_n} import edges resolved to tracked paths (head-only)"
);
Ok(IngestStats::default())
}
}
fn current_head_rev(db: &FactsDb) -> Result<String> {
let sql = "SELECT rev FROM commits ORDER BY date DESC, rowid ASC LIMIT 1";
let mut stmt = db
.conn()
.prepare(sql)
.map_err(|e| CodeLoreError::Analysis(format!("prepare head rev: {e}")))?;
let mut rows = stmt
.query([])
.map_err(|e| CodeLoreError::Analysis(format!("query head rev: {e}")))?;
if let Some(row) = rows
.next()
.map_err(|e| CodeLoreError::Analysis(format!("head rev row: {e}")))?
{
Ok(row
.get::<_, String>(0)
.map_err(|e| CodeLoreError::Analysis(format!("head rev value: {e}")))?)
} else {
Ok(String::new())
}
}
fn query_live_paths(db: &FactsDb) -> Result<Vec<String>> {
let sql = "
WITH latest_per_path AS (
SELECT
c.path,
arg_max(
c.change_type,
ROW(commits.date, -commits.rowid)
) AS change_type
FROM changes c
INNER JOIN commits ON commits.rev = c.rev
GROUP BY c.path
)
SELECT path
FROM latest_per_path
WHERE change_type != 'deleted'
ORDER BY path
";
let mut stmt = db
.conn()
.prepare(sql)
.map_err(|e| CodeLoreError::Analysis(format!("prepare path query: {e}")))?;
stmt.query_map([], |r| r.get::<_, String>(0))
.map_err(|e| CodeLoreError::Analysis(format!("query paths: {e}")))?
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(|e| CodeLoreError::Analysis(format!("collect paths: {e}")))
}
fn future_dated_commit_warning(db: &FactsDb) -> Result<Option<String>> {
let now = crate::analyses::query::wall_clock_utc_literal();
let (count, latest) = db.query_row(
&format!(
"SELECT COUNT(*), CAST(MAX(date) AS TEXT) \
FROM commits WHERE date > TIMESTAMP '{now}'"
),
[],
|r| Ok((r.get::<_, i64>(0)?, r.get::<_, Option<String>>(1)?)),
)?;
Ok((count > 0).then(|| {
format!(
"ingest: {count} commit(s) are dated after the current wall clock \
(latest {}); window anchors are clamped to now — check for a bad \
commit date or contributor clock skew",
latest.as_deref().unwrap_or("unknown")
)
}))
}
pub(crate) fn format_panic_payload(payload: &Box<dyn std::any::Any + Send>) -> String {
let detail = payload
.downcast_ref::<&'static str>()
.map(|s| (*s).to_string())
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "<non-string panic payload>".to_string());
format!("commit walker thread panicked: {detail}")
}
#[cfg(test)]
mod panic_payload_tests {
use super::format_panic_payload;
#[test]
fn extracts_static_str_payload() {
let payload: Box<dyn std::any::Any + Send> = Box::new("walker exploded");
let msg = format_panic_payload(&payload);
assert_eq!(msg, "commit walker thread panicked: walker exploded");
}
#[test]
fn extracts_string_payload() {
let payload: Box<dyn std::any::Any + Send> = Box::new(String::from("formatted reason"));
let msg = format_panic_payload(&payload);
assert_eq!(msg, "commit walker thread panicked: formatted reason");
}
#[test]
fn unknown_payload_falls_through_to_placeholder() {
let payload: Box<dyn std::any::Any + Send> = Box::new(42_u32);
let msg = format_panic_payload(&payload);
assert_eq!(
msg,
"commit walker thread panicked: <non-string panic payload>"
);
}
}
#[cfg(test)]
mod future_date_warning_tests {
use super::future_dated_commit_warning;
use crate::facts::FactsDb;
fn ts_days_ago(days: i64) -> String {
let t = time::OffsetDateTime::now_utc() - time::Duration::days(days);
crate::facts::ingest::consumer::format_timestamp(t)
}
fn seed(db: &FactsDb, rev: &str, date: &str) {
db.execute_batch(&format!(
"INSERT INTO commits (rev, author_email, author_name, committer_email, \
canonical_author, date, committer_date, message, is_merge, parent_count) \
VALUES ('{rev}', 'a@b.c', 'A', 'a@b.c', 'A', TIMESTAMP '{date}', TIMESTAMP '{date}', 'm', false, 1)"
))
.expect("seed commit");
}
#[test]
fn warns_naming_count_and_extent_for_a_future_dated_commit() {
let db = FactsDb::new_in_memory().expect("db");
seed(&db, "r1", &ts_days_ago(5));
seed(&db, "r2", &ts_days_ago(1));
seed(&db, "r3", "2099-01-01 00:00:00");
let msg = future_dated_commit_warning(&db)
.expect("query ok")
.expect("a future-dated commit must warn");
assert!(msg.contains("1 commit"), "names the count: {msg}");
assert!(msg.contains("2099-01-01"), "names the extent: {msg}");
assert!(msg.contains("clamped to now"), "explains the clamp: {msg}");
}
#[test]
fn silent_when_every_commit_predates_now() {
let db = FactsDb::new_in_memory().expect("db");
seed(&db, "r1", &ts_days_ago(5));
seed(&db, "r2", &ts_days_ago(400));
assert!(
future_dated_commit_warning(&db)
.expect("query ok")
.is_none(),
"no future-dated commit ⇒ no warning"
);
}
#[test]
fn silent_on_an_empty_store() {
let db = FactsDb::new_in_memory().expect("db");
assert!(
future_dated_commit_warning(&db)
.expect("query ok")
.is_none()
);
}
}