use rusqlite::{params, Connection};
use crate::collect::ticket::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,
}
impl CorrelationOutcome {
pub fn summary(&self) -> String {
format!(
"{} linked, {} already linked, {} ticket without board item, {} without ticket",
self.linked, self.already_linked, self.no_work_item, self.no_ticket
)
}
}
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)
}
fn ticket_key(ticket_id: Option<&str>, message: &str) -> Option<String> {
ticket_id
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
.or_else(|| extract_ticket_id(message))
}
pub fn correlate_commits(conn: &mut Connection, bus: &ProgressBus) -> Result<CorrelationOutcome> {
let candidates = load_candidates(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) {
None => outcome.no_ticket += 1,
Some(key) => {
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_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", "fixes PROJ-1 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,
};
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"));
}
#[test]
fn ticket_key_prefers_column_then_message() {
assert_eq!(
ticket_key(Some("PROJ-1"), "unrelated PROJ-2"),
Some("PROJ-1".into())
);
assert_eq!(ticket_key(Some(" "), "sees PROJ-2"), Some("PROJ-2".into()));
assert_eq!(ticket_key(None, "sees PROJ-2"), Some("PROJ-2".into()));
assert_eq!(ticket_key(None, "chore: nothing"), None);
}
}