use std::collections::HashMap;
use rusqlite::{params, Connection};
use crate::collect::ticket::{branch_ticket_key, extract_ticket_id};
use crate::core::db::work_items::link_commit_work_item;
use crate::core::errors::{Result, TgaError};
use crate::core::progress::{ProgressBus, ProgressEvent, Stage};
const TARGET: &str = "commit → board item";
const PROGRESS_EVERY: usize = 250;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct CorrelationOutcome {
pub scanned: u64,
pub linked: u64,
pub already_linked: u64,
pub no_ticket: u64,
pub no_work_item: u64,
pub from_branch: u64,
pub from_pr_body: u64,
}
impl CorrelationOutcome {
pub fn summary(&self) -> String {
format!(
"{} linked, {} already linked, {} ticket without board item, \
{} without ticket ({} via branch, {} via PR body)",
self.linked,
self.already_linked,
self.no_work_item,
self.no_ticket,
self.from_branch,
self.from_pr_body
)
}
}
type Candidate = (String, Option<String>, String, bool);
fn load_candidates(conn: &Connection) -> Result<Vec<Candidate>> {
let mut stmt = conn
.prepare(
"SELECT c.sha, c.ticket_id, c.message, \
EXISTS (SELECT 1 FROM commit_work_items w WHERE w.commit_sha = c.sha) \
FROM commits c ORDER BY c.sha",
)
.map_err(TgaError::from)?;
let mapped = stmt
.query_map([], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, Option<String>>(1)?,
r.get::<_, String>(2)?,
r.get::<_, i64>(3)? != 0,
))
})
.map_err(TgaError::from)?;
let mut out = Vec::new();
for row in mapped {
out.push(row.map_err(TgaError::from)?);
}
Ok(out)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum TicketSource {
CommitText,
BranchName,
PrBody,
}
const SOURCE_ORDER: [TicketSource; 3] = [
TicketSource::CommitText,
TicketSource::BranchName,
TicketSource::PrBody,
];
#[derive(Debug, Clone, Default)]
struct PrKeys {
branch: Option<String>,
body: Option<String>,
}
fn load_pr_keys(conn: &Connection) -> Result<HashMap<String, PrKeys>> {
let mut stmt = conn
.prepare(
"SELECT commit_shas, head_ref, body_ticket_id FROM pull_requests \
WHERE (head_ref IS NOT NULL AND head_ref != '') OR body_ticket_id IS NOT NULL",
)
.map_err(TgaError::from)?;
let mapped = stmt
.query_map([], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, Option<String>>(1)?,
r.get::<_, Option<String>>(2)?,
))
})
.map_err(TgaError::from)?;
let mut out: HashMap<String, PrKeys> = HashMap::new();
for row in mapped {
let (shas_json, head_ref, body_key) = row.map_err(TgaError::from)?;
let branch = head_ref.as_deref().and_then(branch_ticket_key);
if branch.is_none() && body_key.is_none() {
continue;
}
let Ok(shas) = serde_json::from_str::<Vec<String>>(&shas_json) else {
continue;
};
for sha in shas {
out.insert(
sha,
PrKeys {
branch: branch.clone(),
body: body_key.clone(),
},
);
}
}
Ok(out)
}
fn ticket_key(
ticket_id: Option<&str>,
message: &str,
pr: Option<&PrKeys>,
) -> Option<(String, TicketSource)> {
SOURCE_ORDER.iter().find_map(|&source| {
let key = match source {
TicketSource::CommitText => ticket_id
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
.or_else(|| extract_ticket_id(message)),
TicketSource::BranchName => pr.and_then(|p| p.branch.clone()),
TicketSource::PrBody => pr.and_then(|p| p.body.clone()),
};
key.map(|k| (k, source))
})
}
pub fn correlate_commits(conn: &mut Connection, bus: &ProgressBus) -> Result<CorrelationOutcome> {
let candidates = load_candidates(conn)?;
let pr_keys = load_pr_keys(conn)?;
let total = candidates.len() as u64;
bus.emit(ProgressEvent::started(
Stage::Correlate,
TARGET,
Some(total),
));
let mut outcome = CorrelationOutcome::default();
let tx = conn.transaction().map_err(TgaError::from)?;
{
let mut lookup = tx
.prepare("SELECT source FROM work_items WHERE id = ?1 ORDER BY source")
.map_err(TgaError::from)?;
for (i, (sha, ticket_id, message, already)) in candidates.iter().enumerate() {
outcome.scanned += 1;
if *already {
outcome.already_linked += 1;
} else {
match ticket_key(ticket_id.as_deref(), message, pr_keys.get(sha)) {
None => outcome.no_ticket += 1,
Some((key, from)) => {
match from {
TicketSource::CommitText => {}
TicketSource::BranchName => outcome.from_branch += 1,
TicketSource::PrBody => outcome.from_pr_body += 1,
}
let sources: Vec<String> = lookup
.query_map(params![key], |r| r.get::<_, String>(0))
.map_err(TgaError::from)?
.collect::<std::result::Result<_, _>>()
.map_err(TgaError::from)?;
if sources.is_empty() {
outcome.no_work_item += 1;
}
for source in &sources {
link_commit_work_item(&tx, sha, &key, source)?;
outcome.linked += 1;
}
}
}
}
if (i + 1).is_multiple_of(PROGRESS_EVERY) {
bus.emit(ProgressEvent::advanced(
Stage::Correlate,
TARGET,
outcome.scanned,
Some(total),
));
}
}
}
tx.commit().map_err(TgaError::from)?;
bus.emit(
ProgressEvent::completed(Stage::Correlate, TARGET, outcome.scanned)
.with_detail(outcome.summary()),
);
Ok(outcome)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::db::correlation::{correlation_counts, correlation_rows, CorrelationFilter};
use crate::core::db::work_items::{upsert_work_item, WorkItemRow};
use crate::core::db::Database;
fn work_item(id: &str, source: &str) -> WorkItemRow {
WorkItemRow {
id: id.into(),
source: source.into(),
title: format!("Item {id}"),
status: "Open".into(),
item_type: "Task".into(),
tags: None,
project: None,
url: None,
raw_json: None,
}
}
fn insert_pr(
conn: &Connection,
pr_number: i64,
sha: &str,
head_ref: Option<&str>,
body_key: Option<&str>,
) {
let shas = serde_json::to_string(&vec![sha]).expect("encode shas");
conn.execute(
"INSERT INTO pull_requests \
(provider, repository, pr_number, title, author, state, created_at, \
commit_shas, head_ref, body_ticket_id) \
VALUES ('github', 'acme/widgets', ?1, 'T', 'ada', 'merged', \
'2026-01-01T00:00:00Z', ?2, ?3, ?4)",
params![pr_number, shas, head_ref.unwrap_or(""), body_key],
)
.expect("insert pr");
}
fn insert_commit(conn: &Connection, sha: &str, message: &str, ticket: Option<&str>) {
conn.execute(
"INSERT INTO commits (sha, author_name, author_email, timestamp, message, \
repository, ticket_id) \
VALUES (?1, 'A', 'a@x', '2026-01-01T00:00:00Z', ?2, 'repo', ?3)",
params![sha, message, ticket],
)
.expect("insert commit");
}
#[test]
fn links_matching_ticket_keys() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "PROJ-1 x", Some("PROJ-1"));
insert_commit(
db.connection(),
"bbb",
"PROJ-9 no such item",
Some("PROJ-9"),
);
insert_commit(db.connection(), "ccc", "chore: nothing here", None);
insert_commit(db.connection(), "ddd", "PROJ-1 fixed again", None);
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("upsert");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.scanned, 4);
assert_eq!(out.linked, 2, "aaa and ddd both resolve PROJ-1");
assert_eq!(out.no_work_item, 1, "PROJ-9 has no board item");
assert_eq!(out.no_ticket, 1);
assert_eq!(out.already_linked, 0);
assert_eq!(
correlation_counts(db.connection()).expect("counts").linked,
2
);
}
#[test]
fn produces_full_result_with_no_credentials_configured() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "ENG-7 ship it", Some("ENG-7"));
insert_commit(db.connection(), "bbb", "AB#42 azdo work", Some("AB#42"));
insert_commit(db.connection(), "ccc", "chore: no ticket", None);
upsert_work_item(db.connection(), &work_item("ENG-7", "linear")).expect("linear");
upsert_work_item(db.connection(), &work_item("AB#42", "azdo")).expect("azdo");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 2);
assert_eq!(out.no_ticket, 1);
let counts = correlation_counts(db.connection()).expect("counts");
assert_eq!(counts.commits, 3);
assert_eq!(counts.linked, 2);
assert_eq!(counts.unticketed, 1);
assert_eq!(counts.work_items_linked, 2);
let rows =
correlation_rows(db.connection(), CorrelationFilter::Unlinked, 10).expect("rows");
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].sha, "ccc", "the gap is named, not guessed at");
}
#[test]
fn is_idempotent() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "PROJ-1 x", Some("PROJ-1"));
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("upsert");
let first = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("1st");
assert_eq!(first.linked, 1);
let second = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("2nd");
assert_eq!(second.linked, 0);
assert_eq!(second.already_linked, 1);
assert_eq!(
correlation_counts(db.connection()).expect("counts").linked,
1
);
}
#[test]
fn links_the_same_key_across_sources() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "PROJ-1 x", Some("PROJ-1"));
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("jira");
upsert_work_item(db.connection(), &work_item("PROJ-1", "linear")).expect("linear");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 2);
}
#[test]
fn never_invents_a_work_item() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "PROJ-1 x", Some("PROJ-1"));
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 0);
assert_eq!(out.no_work_item, 1);
assert_eq!(
correlation_counts(db.connection())
.expect("counts")
.work_items,
0
);
}
#[test]
fn emits_started_and_completed() {
let mut db = Database::open_in_memory().expect("open");
for i in 0..3 {
insert_commit(db.connection(), &format!("sha{i}"), "chore: x", None);
}
let bus = ProgressBus::bounded(64);
correlate_commits(db.connection_mut(), &bus).expect("run");
let events = bus.drain();
assert_eq!(events.len(), 2, "one started, one completed (3 < 250)");
assert_eq!(events[0].stage, Stage::Correlate);
assert_eq!(events[0].target, TARGET);
assert_eq!(events[0].total, Some(3));
assert!(!events[0].is_terminal());
assert!(events[1].is_terminal());
assert_eq!(events[1].done, 3);
}
#[test]
fn disabled_bus_produces_identical_outcome() {
let build = || {
let db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "PROJ-1 x", Some("PROJ-1"));
insert_commit(db.connection(), "bbb", "chore: y", None);
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("upsert");
db
};
let mut off = build();
let mut on = build();
let a = correlate_commits(off.connection_mut(), &ProgressBus::disabled()).expect("off");
let b = correlate_commits(on.connection_mut(), &ProgressBus::bounded(8)).expect("on");
assert_eq!(a, b);
}
#[test]
fn summary_names_every_disposition() {
let out = CorrelationOutcome {
scanned: 4,
linked: 1,
already_linked: 1,
no_ticket: 1,
no_work_item: 1,
from_branch: 2,
from_pr_body: 3,
};
let s = out.summary();
assert!(s.contains("1 linked"));
assert!(s.contains("1 already linked"));
assert!(s.contains("1 ticket without board item"));
assert!(s.contains("1 without ticket"));
assert!(s.contains("2 via branch"));
assert!(s.contains("3 via PR body"));
}
#[test]
fn ticket_key_prefers_column_then_message() {
assert_eq!(
ticket_key(Some("PROJ-1"), "PROJ-2 unrelated", None),
Some(("PROJ-1".into(), TicketSource::CommitText))
);
assert_eq!(
ticket_key(Some(" "), "PROJ-2 seen", None),
Some(("PROJ-2".into(), TicketSource::CommitText))
);
assert_eq!(
ticket_key(None, "PROJ-2 seen", None),
Some(("PROJ-2".into(), TicketSource::CommitText))
);
assert_eq!(ticket_key(None, "chore: nothing", None), None);
}
#[test]
fn adr_citation_does_not_become_a_phantom_coverage_gap() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(
db.connection(),
"352fe5d6",
"fix(relay): verify the HMAC once\n\nPer ADR-0034 the relay spools durably.\nCloses #5089\n",
None,
);
upsert_work_item(db.connection(), &work_item("#5089", "github")).expect("upsert");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(
out.linked, 1,
"resolves the issue the commit actually cites"
);
assert_eq!(out.no_work_item, 0, "ADR-0034 is not a phantom gap");
assert_eq!(out.no_ticket, 0);
}
#[test]
fn branch_name_supplies_a_key_the_subject_lacks() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "chore: tidy up", None);
insert_pr(db.connection(), 1, "aaa", Some("feature/PROJ-1-tidy"), None);
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("upsert");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 1, "the branch name carried the key");
assert_eq!(out.from_branch, 1);
assert_eq!(out.from_pr_body, 0);
assert_eq!(out.no_ticket, 0);
}
#[test]
fn pr_body_supplies_a_key_the_subject_lacks() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "chore: tidy up", None);
insert_pr(db.connection(), 1, "aaa", None, Some("#5089"));
upsert_work_item(db.connection(), &work_item("#5089", "github")).expect("upsert");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 1);
assert_eq!(out.from_pr_body, 1);
assert_eq!(out.from_branch, 0);
}
#[test]
fn precedence_prefers_commit_text_then_branch_then_body() {
let pr = PrKeys {
branch: Some("BRANCH-2".into()),
body: Some("#3".into()),
};
assert_eq!(
ticket_key(Some("SUBJ-1"), "SUBJ-1 x", Some(&pr)),
Some(("SUBJ-1".into(), TicketSource::CommitText)),
"the commit's own text outranks both pull-request sources"
);
assert_eq!(
ticket_key(None, "chore: nothing", Some(&pr)),
Some(("BRANCH-2".into(), TicketSource::BranchName)),
"a deliberate branch identifier outranks PR prose"
);
let body_only = PrKeys {
branch: None,
body: Some("#3".into()),
};
assert_eq!(
ticket_key(None, "chore: nothing", Some(&body_only)),
Some(("#3".into(), TicketSource::PrBody))
);
assert_eq!(ticket_key(None, "chore: nothing", None), None);
}
#[test]
fn precedence_order_matches_the_declared_variant_order() {
let mut sorted = SOURCE_ORDER;
sorted.sort();
assert_eq!(SOURCE_ORDER, sorted);
assert!(TicketSource::CommitText < TicketSource::BranchName);
assert!(TicketSource::BranchName < TicketSource::PrBody);
}
#[test]
fn branch_and_body_noise_never_becomes_a_phantom_coverage_gap() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "chore: tidy up", None);
insert_pr(
db.connection(),
1,
"aaa",
Some("fix/ADR-0029-followup"),
None,
);
insert_commit(db.connection(), "bbb", "chore: tidy more", None);
insert_pr(db.connection(), 2, "bbb", Some("fix/5734-slug"), None);
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.no_ticket, 2, "neither commit gains a key");
assert_eq!(out.from_branch, 0);
assert_eq!(out.no_work_item, 0, "no phantom gap is reported");
}
#[test]
fn summary_reports_a_zero_harvest_rather_than_omitting_it() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "chore: nothing", None);
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.from_branch, 0);
assert_eq!(out.from_pr_body, 0);
let s = out.summary();
assert!(s.contains("0 via branch"), "summary was: {s}");
assert!(s.contains("0 via PR body"), "summary was: {s}");
}
#[test]
fn a_malformed_commit_shas_column_skips_one_pr_not_the_pass() {
let mut db = Database::open_in_memory().expect("open");
insert_commit(db.connection(), "aaa", "chore: tidy", None);
insert_commit(db.connection(), "bbb", "chore: tidy", None);
db.connection()
.execute(
"INSERT INTO pull_requests \
(provider, repository, pr_number, title, author, state, created_at, \
commit_shas, head_ref) \
VALUES ('github','acme/w',1,'T','ada','merged','2026-01-01T00:00:00Z', \
'not json', 'feature/PROJ-1-x')",
[],
)
.expect("insert malformed pr");
insert_pr(db.connection(), 2, "bbb", Some("feature/PROJ-1-y"), None);
upsert_work_item(db.connection(), &work_item("PROJ-1", "jira")).expect("upsert");
let out = correlate_commits(db.connection_mut(), &ProgressBus::disabled()).expect("run");
assert_eq!(out.linked, 1, "the well-formed PR still correlates");
assert_eq!(out.from_branch, 1);
assert_eq!(out.no_ticket, 1, "the malformed PR's commit gains nothing");
}
}