use chrono::{DateTime, NaiveDate, TimeZone, Utc};
use tracing::{info, warn};
use crate::collect::azdo::AzureDevOpsClient;
use crate::collect::bitbucket::BitbucketClient;
use crate::collect::errors::Result;
use crate::collect::fault::CollectionFault;
use crate::collect::git::walk_state;
use crate::collect::git::GitCollector;
use crate::collect::github::budget::RunBudget;
use crate::collect::github::GitHubClient;
use crate::collect::identity::IdentityResolver;
use crate::collect::linear_pipeline;
use crate::collect::notify;
use crate::collect::pr_provider::PrProvider;
use crate::collect::reclassify::reclassify_stale;
use crate::collect::weeks::{clamp_week_to_range, weeks_in_range};
use crate::collect::work_item_pipeline;
use crate::core::config::Config;
use crate::core::db::{self, Database};
use crate::core::models::PullRequest;
use crate::core::progress::{ProgressBus, ProgressEvent, Stage};
#[derive(Debug, Clone)]
pub enum FetchOutcome {
Success {
remote: String,
},
Failed {
remote: String,
error: String,
},
Skipped {
reason: String,
},
}
#[derive(Debug, Clone)]
pub struct PerRepoFetch {
pub repo: String,
pub outcome: FetchOutcome,
}
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct CollectionStats {
pub commits_collected: usize,
pub authors_resolved: usize,
pub prs_fetched: usize,
pub reviewers_fetched: usize,
pub linear_issues_fetched: usize,
pub weeks_collected: usize,
pub weeks_skipped: usize,
pub repos_skipped: usize,
pub errors: Vec<CollectionFault>,
pub reachability_rows: usize,
pub fetch_outcomes: Vec<PerRepoFetch>,
}
pub struct CollectionPipeline {
config: Config,
force: bool,
no_fetch: bool,
force_refresh_prs: bool,
skip_tag_reachability: bool,
head_only: bool,
branches: Vec<String>,
strict_fetch: bool,
verbose_fetch: bool,
progress: ProgressBus,
github_budget: RunBudget,
}
impl CollectionPipeline {
pub fn new(config: Config) -> Self {
Self {
config,
force: false,
no_fetch: false,
force_refresh_prs: false,
skip_tag_reachability: false,
head_only: false,
branches: Vec::new(),
strict_fetch: false,
verbose_fetch: false,
progress: ProgressBus::disabled(),
github_budget: RunBudget::new(),
}
}
#[must_use]
pub fn with_progress(mut self, progress: ProgressBus) -> Self {
self.progress = progress;
self
}
pub fn with_force(mut self, force: bool) -> Self {
self.force = force;
self
}
pub fn with_no_fetch(mut self, no_fetch: bool) -> Self {
self.no_fetch = no_fetch;
self
}
pub fn with_skip_tag_reachability(mut self, skip: bool) -> Self {
self.skip_tag_reachability = skip;
self
}
pub fn with_head_only(mut self, head_only: bool) -> Self {
self.head_only = head_only;
self
}
pub fn with_branches(mut self, branches: Vec<String>) -> Self {
self.branches = branches;
self
}
pub fn with_strict_fetch(mut self, strict: bool) -> Self {
self.strict_fetch = strict;
self
}
pub fn with_verbose_fetch(mut self, verbose: bool) -> Self {
self.verbose_fetch = verbose;
self
}
pub fn strict_fetch(&self) -> bool {
self.strict_fetch
}
pub fn verbose_fetch(&self) -> bool {
self.verbose_fetch
}
pub fn with_force_refresh_prs(mut self, force_refresh_prs: bool) -> Self {
self.force_refresh_prs = force_refresh_prs;
self
}
pub fn config(&self) -> &Config {
&self.config
}
fn reclassify_stale_commits(&self, db: &mut Database) -> Result<()> {
let stats = reclassify_stale(db)?;
if stats.stamped > 0 {
notify::warning(
&self.progress,
"reclassify",
&format!(
"Re-classified {} commit(s) stored by an older AI detector; \
{} verdict(s) changed.",
stats.stamped, stats.changed
),
);
}
Ok(())
}
pub async fn run(&self, db: &mut Database) -> Result<CollectionStats> {
let mut stats = CollectionStats::default();
self.reclassify_stale_commits(db)?;
let resolver = IdentityResolver::from_config(&self.config);
for repo_cfg in self.config.repositories.iter() {
let repo_label = repo_cfg
.name
.clone()
.unwrap_or_else(|| repo_cfg.path.display().to_string());
self.progress.emit(ProgressEvent::started(
Stage::Collect,
repo_label.clone(),
None,
));
let effective_head_only = self.head_only || repo_cfg.head_only;
let pre_fetch_collector = match GitCollector::new(repo_cfg) {
Ok(c) => c
.no_fetch(self.no_fetch)
.with_head_only(effective_head_only)
.with_explicit_branches(self.branches.clone()),
Err(e) => {
let msg = format!("failed to open repo {}: {e}", repo_cfg.path.display());
warn!("{msg}");
self.progress.emit(ProgressEvent::failed(
Stage::Collect,
repo_label.clone(),
msg.clone(),
));
stats.fail_stage(msg);
continue;
}
};
let fetch_result = pre_fetch_collector.perform_fetch();
stats.fetch_outcomes.push(fetch_result);
let collector = match GitCollector::new(repo_cfg) {
Ok(c) => c
.no_fetch(true)
.with_head_only(effective_head_only)
.with_explicit_branches(self.branches.clone())
.with_progress(self.progress.clone()),
Err(e) => {
let msg = format!("failed to open repo {}: {e}", repo_cfg.path.display());
warn!("{msg}");
self.progress.emit(ProgressEvent::failed(
Stage::Collect,
repo_label.clone(),
msg.clone(),
));
stats.fail_stage(msg);
continue;
}
};
let before = stats.commits_collected as u64;
let errors_before = stats.errors.len();
let failures = self.collect_repo_by_week(db, &collector, &mut stats);
let collected = stats.commits_collected as u64 - before;
self.progress.emit(if failures == 0 {
ProgressEvent::completed(Stage::Collect, repo_label, collected)
} else {
let first = stats
.errors
.get(errors_before)
.map_or("", |f| f.message.as_str());
ProgressEvent::failed(
Stage::Collect,
repo_label,
format!("{failures} error(s), {collected} commit(s) collected; {first}"),
)
});
}
if !self.skip_tag_reachability {
self.run_reachability_scan(db, &mut stats);
} else {
info!("skipping tag/release-branch reachability scan (--skip-tag-reachability)");
}
stats.authors_resolved = self.upsert_observed_authors(db, &resolver)?;
if let Ok(unresolved) = count_unresolved_commits(db) {
if unresolved > 0 {
let msg = format!(
"WARNING: {unresolved} commits have unresolved author identities and may \
inflate developer counts. Run `tga aliases list` to review, or extend \
`developer_aliases` in the config to map missing identities."
);
warn!("{msg}");
notify::warning(&self.progress, "collect", &msg);
}
}
self.fetch_and_store_prs(db, &mut stats).await;
if let Some(azdo_cfg) = self.config.azure_devops_config() {
let client = AzureDevOpsClient::new(azdo_cfg.clone());
match client.test_connection().await {
Ok(info) => info!(
user = info.user_name.as_deref().unwrap_or("?"),
org = %info.organization_url,
"Azure DevOps connection verified",
),
Err(e) => {
warn!("Azure DevOps connection failed (non-fatal): {e}");
}
}
if azdo_cfg.fetch_prs {
match self.fetch_and_persist_azdo_prs(db, azdo_cfg).await {
Ok(n) => {
info!(prs = n, "stored ADO pull requests");
stats.prs_fetched += n;
}
Err(e) => {
stats.fail_stage(format!("ADO PR fetch failed: {e}"));
}
}
}
}
linear_pipeline::fetch_and_store_linear_issues(db, &self.config, &mut stats).await;
work_item_pipeline::fetch_and_persist_work_items(db, &self.config, &mut stats).await;
Ok(stats)
}
fn run_reachability_scan(&self, db: &mut Database, stats: &mut CollectionStats) {
use crate::collect::git::reachability::scan_and_persist;
use crate::core::config::expand_path;
let cfg = &self.config.reachability;
if !cfg.track_tags && !cfg.track_release_branches {
info!("reachability tracking disabled by config (track_tags=false, track_release_branches=false)");
return;
}
let conn = db.connection();
for repo_cfg in &self.config.repositories {
let path = expand_path(&repo_cfg.path);
let name = repo_cfg
.name
.clone()
.or_else(|| {
path.file_name()
.and_then(|s| s.to_str())
.map(|s| s.to_string())
})
.unwrap_or_else(|| path.display().to_string());
info!(repo = %name, "running reachability scan");
match scan_and_persist(&path, conn, cfg, Some(&name)) {
Ok(r) => {
info!(
repo = %name,
rows = r.rows_upserted,
default_branch = r.default_branch_commits,
tagged = r.tagged_commits,
release_branch = r.release_branch_commits,
"reachability scan complete"
);
stats.reachability_rows += r.rows_upserted;
}
Err(e) => {
let msg = format!("reachability scan failed for {name}: {e}");
warn!("{msg}");
stats.fail_stage(msg);
}
}
}
}
fn build_pr_providers(
&self,
stats: &mut CollectionStats,
org_discovered: &[(String, String)],
workspace_discovered: &crate::collect::bitbucket::WorkspaceDiscovery,
) -> Vec<Box<dyn PrProvider + Send + Sync>> {
let mut providers: Vec<Box<dyn PrProvider + Send + Sync>> = Vec::new();
if let Some(gh_cfg) = &self.config.github {
if gh_cfg.fetch_prs {
let repos =
crate::collect::github::org_discovery::resolve_github_repos_with_discovered(
gh_cfg,
&self.config.repositories,
org_discovered,
);
if repos.is_empty() {
info!(
"GitHub PR fetch skipped: no github.repo, no per-repo org, \
no github.org/orgs resolvable from repositories[] or org discovery"
);
} else if gh_cfg.token.is_none() && std::env::var("GITHUB_TOKEN").ok().is_none() {
let msg = "GitHub PR fetch is enabled (github.fetch_prs=true) but \
no token is configured. Set `github.token` or the \
GITHUB_TOKEN env var to a PAT with `repo` scope (public \
repos only need `public_repo`); without it, GitHub \
rate-limits to 60 requests/hour and most PRs will be \
missed.";
warn!("{msg}");
notify::warning(&self.progress, "collect", &format!("warning: {msg}"));
info!(
repo_count = repos.len(),
"GitHub PR fetcher will scan {} repo(s) anonymously",
repos.len()
);
match GitHubClient::new_for_prs(gh_cfg, repos)
.map(|c| c.with_run_budget(&self.github_budget))
{
Ok(gh) => providers.push(Box::new(gh)),
Err(e) => stats.fail_stage(format!("GitHub client init failed: {e}")),
}
} else {
info!(
repo_count = repos.len(),
"GitHub PR fetcher will scan {} repo(s)",
repos.len()
);
match GitHubClient::new_for_prs(gh_cfg, repos)
.map(|c| c.with_run_budget(&self.github_budget))
{
Ok(gh) => providers.push(Box::new(gh)),
Err(e) => stats.fail_stage(format!("GitHub client init failed: {e}")),
}
}
} else {
info!(
"GitHub PR fetch disabled (github.fetch_prs=false). Set \
`github.fetch_prs: true` in your config to populate the \
pull_requests table."
);
}
} else if has_github_like_repos(&self.config.repositories) {
let msg = "Repositories look like GitHub clones, but no `github:` config \
block is present. To populate the `pull_requests` table, add:\n\
\n\
github:\n \
token: \"${GITHUB_TOKEN}\" # PAT with `repo` scope\n \
fetch_prs: true\n \
repo: \"owner/name\" # OR `org: \"owner\"` for org-wide\n";
tracing::info!("{msg}");
}
if let Some(bb_cfg) = &self.config.bitbucket {
if bb_cfg.fetch_prs {
let repos = crate::collect::bitbucket::resolve_bitbucket_repos(
bb_cfg,
&workspace_discovered.repos,
);
if repos.is_empty() {
info!(
"Bitbucket PR fetch skipped: no bitbucket.workspace/repo_slug pair \
and bitbucket.workspaces discovery returned no repositories"
);
} else {
info!(
repo_count = repos.len(),
"Bitbucket PR fetcher will scan {} repo(s)",
repos.len()
);
match BitbucketClient::new_for_repos(bb_cfg, repos)
.map(|c| c.with_notices(workspace_discovered.notices.clone()))
{
Ok(bb) => providers.push(Box::new(bb)),
Err(e) => stats.fail_stage(format!("Bitbucket client init failed: {e}")),
}
}
}
}
providers
}
async fn fetch_and_store_prs(&self, db: &mut Database, stats: &mut CollectionStats) {
let org_discovered = if let Some(gh_cfg) = &self.config.github {
if gh_cfg.fetch_prs && (!gh_cfg.orgs.is_empty() || gh_cfg.org.is_some()) {
super::github_pipeline::run_github_org_discovery(gh_cfg, &self.github_budget).await
} else {
Vec::new()
}
} else {
Vec::new()
};
let workspace_discovered = match &self.config.bitbucket {
Some(bb_cfg) if bb_cfg.fetch_prs => {
crate::collect::bitbucket::run_workspace_discovery(bb_cfg).await
}
_ => crate::collect::bitbucket::WorkspaceDiscovery::default(),
};
let providers = self.build_pr_providers(stats, &org_discovered, &workspace_discovered);
if providers.is_empty() {
return;
}
let mut set: tokio::task::JoinSet<(String, Result<Vec<PullRequest>>)> =
tokio::task::JoinSet::new();
let providers: Vec<std::sync::Arc<dyn PrProvider + Send + Sync>> =
providers.into_iter().map(std::sync::Arc::from).collect();
for p in &providers {
let p = std::sync::Arc::clone(p);
let name = p.name().to_string();
set.spawn(async move {
let result = p.fetch_pull_requests().await;
(name, result)
});
}
super::pr_pipeline::drain_and_store_pull_requests(set, &providers, db, stats).await;
if let Some(gh_cfg) = &self.config.github {
if gh_cfg.fetch_prs && gh_cfg.fetch_pr_reviews {
super::github_pipeline::fetch_and_store_github_reviewers(
db,
gh_cfg,
self.force_refresh_prs,
stats,
&self.github_budget,
)
.await;
}
}
}
fn collect_repo_by_week(
&self,
db: &mut Database,
collector: &GitCollector,
stats: &mut CollectionStats,
) -> usize {
let errors_before = stats.errors.len();
let repo_name = collector.name().to_string();
let (from, to) = match (collector.since(), collector.until()) {
(Some(s), Some(u)) => (s.date_naive(), u.date_naive()),
(Some(s), None) => (s.date_naive(), Utc::now().date_naive()),
(None, Some(u)) => {
warn!(
repo = %repo_name,
"until_date set without since_date — collecting full git history. \
Use --weeks N or set analysis.since_date in config to limit scope."
);
notify::warning(
&self.progress,
&repo_name,
&format!(
"warning: [{repo_name}] no since_date / --weeks — collecting FULL git history. \
Set analysis.since_date or pass --weeks N to limit scope."
),
);
match collector.collect_window(db, None, Some(u)) {
Ok(n) => {
info!(repo = %repo_name, commits = n, "extracted (until-only)");
stats.commits_collected += n;
}
Err(e) => {
let msg = format!("collection failed for {repo_name}: {e}");
warn!("{msg}");
stats.fail_stage(msg);
}
}
return stats.errors.len() - errors_before;
}
(None, None) => {
warn!(
repo = %repo_name,
"no since_date or --weeks flag set — collecting full git history. \
Use --weeks N or set analysis.since_date in config to limit scope."
);
notify::warning(
&self.progress,
&repo_name,
&format!(
"warning: [{repo_name}] no since_date / --weeks — collecting FULL git history. \
Set analysis.since_date or pass --weeks N to limit scope."
),
);
self.collect_unbounded(db, collector, stats);
return stats.errors.len() - errors_before;
}
};
for week in weeks_in_range(from, to) {
let (year, week_no, _, _) = week;
if !self.force {
match db::is_week_collected(db, &repo_name, year, week_no) {
Ok(true) => {
info!("Skipping {repo_name} W{week_no} {year} — already collected");
notify::progress(
&self.progress,
&repo_name,
&format!(
"Skipped W{week_no:02} {year}: already collected \
(use --force to re-collect) [{repo_name}]"
),
);
stats.weeks_skipped += 1;
continue;
}
Ok(false) => {}
Err(e) => {
let msg = format!(
"collection_runs lookup failed for {repo_name} W{week_no} {year}: {e}"
);
warn!("{msg}");
stats.skip_item(msg);
continue;
}
}
}
let (win_start, win_end) = clamp_week_to_range(week, from, to);
let since_ts = naive_date_start_utc(win_start);
let until_ts = naive_date_end_utc(win_end);
match collector.collect_window(db, Some(since_ts), Some(until_ts)) {
Ok(n) => {
info!(
repo = %repo_name,
year,
week = week_no,
commits = n,
"extracted week"
);
let line = format!("Collected W{week_no:02} {year}: {n} commits [{repo_name}]");
notify::progress(&self.progress, &repo_name, &line);
stats.commits_collected += n;
stats.weeks_collected += 1;
let repo_count = self.config.repositories.len();
if let Err(e) =
db::record_collection_run(db, &repo_name, year, week_no, n, repo_count)
{
let msg = format!(
"failed to record collection_run for {repo_name} W{week_no} {year}: {e}"
);
warn!("{msg}");
stats.skip_item(msg);
}
}
Err(e) => {
let msg = format!("collection failed for {repo_name} W{week_no} {year}: {e}");
warn!("{msg}");
stats.skip_item(msg);
}
}
}
stats.errors.len() - errors_before
}
pub(crate) fn collect_unbounded(
&self,
db: &mut Database,
collector: &GitCollector,
stats: &mut CollectionStats,
) {
let repo_name = collector.name().to_string();
let scope = collector.walk_scope();
let tips = match collector.walk_tips() {
Ok(t) => Some(t),
Err(e) => {
warn!(repo = %repo_name, error = %e, "could not read repository tips; walking in full");
None
}
};
let plan = match (&tips, self.force) {
(_, true) => walk_state::WalkPlan::Full {
reason: walk_state::FullWalkReason::Forced,
},
(None, _) => walk_state::WalkPlan::Full {
reason: walk_state::FullWalkReason::NeverWalked,
},
(Some(t), false) => {
let recorded = walk_state::load(db.connection(), &repo_name).unwrap_or_else(|e| {
warn!(repo = %repo_name, error = %e, "could not read walk state; walking in full");
None
});
let reachable = recorded
.as_ref()
.is_some_and(|s| collector.base_is_reachable(&s.head_sha));
walk_state::plan(recorded.as_ref(), t, &scope, reachable)
}
};
let hide = match &plan {
walk_state::WalkPlan::Skip => {
let line = format!(
"Skipped full history: extract db already current at {head} \
(use --force to re-walk) [{repo_name}]",
head = tips.as_ref().map(|t| t.head_sha.as_str()).unwrap_or("")
);
info!(repo = %repo_name, "{line}");
notify::progress(&self.progress, &repo_name, &line);
stats.repos_skipped += 1;
return;
}
walk_state::WalkPlan::Incremental { base_sha } => match git2::Oid::from_str(base_sha) {
Ok(oid) => Some(oid),
Err(e) => {
warn!(repo = %repo_name, error = %e, "recorded walk base is malformed; walking in full");
None
}
},
walk_state::WalkPlan::Full { reason } => {
info!(
repo = %repo_name,
"collecting FULL git history: {}",
reason.as_str()
);
None
}
};
if let Err(e) = walk_state::mark_in_flight(db.connection(), &repo_name) {
warn!(repo = %repo_name, error = %e, "could not mark walk in flight");
}
match collector.collect_window_hiding(db, None, None, hide) {
Ok(n) => {
info!(repo = %repo_name, commits = n, ?hide, "extracted (unbounded)");
stats.commits_collected += n;
if let Some(t) = &tips {
if let Err(e) =
walk_state::record_complete(db.connection(), &repo_name, t, &scope)
{
warn!(repo = %repo_name, error = %e, "could not record completed walk");
}
}
}
Err(e) => {
let msg = format!("collection failed for {repo_name}: {e}");
warn!("{msg}");
stats.fail_stage(msg);
}
}
}
fn upsert_observed_authors(
&self,
db: &mut Database,
resolver: &IdentityResolver,
) -> Result<usize> {
let pairs: Vec<(String, String)> = {
let conn = db.connection();
let mut stmt = conn.prepare(
"SELECT DISTINCT author_name, author_email FROM commits WHERE author_id IS NULL",
)?;
let rows = stmt.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})?;
let mut out = Vec::new();
for r in rows {
out.push(r?);
}
out
};
let mut count = 0usize;
for (name, email) in pairs {
let author_id = resolver.upsert_author(db, &name, &email)?;
db.connection().execute(
"UPDATE commits SET author_id = ?1 \
WHERE author_id IS NULL AND author_name = ?2 AND author_email = ?3",
rusqlite::params![author_id, name, email],
)?;
count += 1;
}
Ok(count)
}
async fn fetch_and_persist_azdo_prs(
&self,
db: &mut Database,
azdo_cfg: &crate::core::config::AzureDevOpsConfig,
) -> Result<usize> {
use crate::collect::azdo::AdoPrFetcher;
let messages: Vec<String> = {
let conn = db.connection();
let mut stmt = conn.prepare("SELECT message FROM commits")?;
let rows = stmt.query_map([], |row| row.get::<_, String>(0))?;
let mut out = Vec::new();
for r in rows {
out.push(r?);
}
out
};
let fetcher = match AdoPrFetcher::new(azdo_cfg.clone()) {
Ok(f) => f,
Err(e) => {
warn!("ADO PR fetcher init failed: {e}");
return Ok(0);
}
};
let conn = db.connection();
let stored = fetcher
.run_with_options(
conn,
messages.iter().map(String::as_str),
self.force_refresh_prs,
)
.await?;
Ok(stored)
}
}
pub(crate) fn has_github_like_repos(
repositories: &[crate::core::config::RepositoryConfig],
) -> bool {
for repo_cfg in repositories {
let Ok(repo) = git2::Repository::open(&repo_cfg.path) else {
continue;
};
let Ok(remote) = repo.find_remote("origin") else {
continue;
};
let Some(url) = remote.url() else {
continue;
};
if url.contains("github.com") {
return true;
}
}
false
}
fn naive_date_start_utc(d: NaiveDate) -> DateTime<Utc> {
let ndt = d
.and_hms_opt(0, 0, 0)
.expect("00:00:00 is always a valid time");
Utc.from_utc_datetime(&ndt)
}
fn count_unresolved_commits(db: &Database) -> Result<usize> {
let n: i64 = db
.connection()
.query_row(
"SELECT COUNT(*) FROM commits WHERE author_id IS NULL",
[],
|r| r.get(0),
)
.map_err(crate::core::TgaError::from)?;
Ok(n as usize)
}
fn naive_date_end_utc(d: NaiveDate) -> DateTime<Utc> {
let ndt = d
.and_hms_opt(23, 59, 59)
.expect("23:59:59 is always a valid time");
Utc.from_utc_datetime(&ndt)
}