use rusqlite::params;
use crate::collect::ai_marker_config::MarkerScope;
use crate::collect::ai_markers::{detect, CommitSignals, Detection, DETECTOR_VERSION};
use crate::collect::errors::Result;
use crate::core::db::Database;
pub const RECLASSIFY_BATCH: usize = 1_000;
type ScannedRow = (i64, String, Option<String>, String, String, Option<String>);
type PendingWrite = (
i64,
i64,
Option<String>,
&'static str,
Option<&'static str>,
bool,
);
const SCAN_SQL: &str =
"SELECT id, message, ai_tool, agentic_mode, author_email, ai_detection_method FROM commits \
INDEXED BY idx_commits_ai_detector_version \
WHERE ai_detector_version < ?1 \
ORDER BY ai_detector_version, id LIMIT ?2";
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct ReclassifyStats {
pub stamped: usize,
pub changed: usize,
}
pub fn reclassify_stale(db: &mut Database) -> Result<ReclassifyStats> {
reclassify_stale_with(db, RECLASSIFY_BATCH, detect)
}
pub fn reclassify_stale_with<F>(
db: &mut Database,
batch: usize,
mut detector: F,
) -> Result<ReclassifyStats>
where
F: FnMut(&CommitSignals<'_>) -> Detection,
{
let batch = batch.max(1);
let mut total = ReclassifyStats::default();
loop {
let pass = reclassify_batch_with(db, batch, &mut detector)?;
if pass.stamped == 0 {
return Ok(total);
}
total.stamped += pass.stamped;
total.changed += pass.changed;
}
}
pub fn reclassify_batch_with<F>(
db: &mut Database,
batch: usize,
detector: &mut F,
) -> Result<ReclassifyStats>
where
F: FnMut(&CommitSignals<'_>) -> Detection,
{
let mut pending: Vec<PendingWrite> = Vec::new();
{
let conn = db.connection();
let mut stmt = conn.prepare(SCAN_SQL)?;
let rows: Vec<ScannedRow> = stmt
.query_map(params![DETECTOR_VERSION, batch as i64], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, Option<String>>(3)?
.unwrap_or_else(|| "none".to_string()),
row.get::<_, Option<String>>(4)?.unwrap_or_default(),
row.get::<_, Option<String>>(5)?,
))
})?
.collect::<std::result::Result<_, _>>()?;
for (id, message, stored_tool, stored_mode, author_email, stored_method) in rows {
let detection = detector(&CommitSignals {
message: &message,
author_email: &author_email,
committer_email: "",
});
let tool = detection.tool;
let mode = detection.mode.as_str();
let method = detection.method.map(MarkerScope::as_str);
let changed = tool != stored_tool.as_deref()
|| mode != stored_mode
|| method != stored_method.as_deref();
let is_ai = i64::from(tool.is_some());
pending.push((id, is_ai, tool.map(str::to_string), mode, method, changed));
}
}
if pending.is_empty() {
return Ok(ReclassifyStats::default());
}
let stats = ReclassifyStats {
stamped: pending.len(),
changed: pending.iter().filter(|(.., changed)| *changed).count(),
};
let conn = db.connection_mut();
let tx = conn.transaction()?;
{
let mut up = tx.prepare(
"UPDATE commits SET is_ai_assisted = ?1, ai_tool = ?2, agentic_mode = ?3, \
ai_detector_version = ?4, ai_detection_method = ?5 WHERE id = ?6",
)?;
for (id, is_ai, tool, mode, method, _) in &pending {
up.execute(params![is_ai, tool, mode, DETECTOR_VERSION, method, id])?;
}
}
tx.commit()?;
Ok(stats)
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
fn insert_commit(db: &Database, sha: &str, message: &str, detector_version: i64) {
db.connection()
.execute(
"INSERT INTO commits \
(sha, author_name, author_email, timestamp, message, repository, \
is_ai_assisted, ai_tool, agentic_mode, ai_detector_version) \
VALUES (?1, 'Ada', 'ada@example.com', '2026-01-01T00:00:00Z', ?2, \
'testrepo', 0, NULL, 'none', ?3)",
params![sha, message, detector_version],
)
.expect("insert commit");
}
fn verdict(db: &Database, sha: &str) -> (i64, Option<String>, String, i64) {
db.connection()
.query_row(
"SELECT is_ai_assisted, ai_tool, agentic_mode, ai_detector_version \
FROM commits WHERE sha = ?1",
params![sha],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
)
.expect("read verdict")
}
const TRAILER: &str = "feat: a thing\n\nCo-Authored-By: Claude <noreply@anthropic.com>\n";
#[test]
fn a_stale_trailer_commit_is_repaired() {
let mut db = Database::open_in_memory().expect("open db");
insert_commit(&db, "stale_trailer", TRAILER, 0);
let stats = reclassify_stale(&mut db).expect("reclassify");
assert_eq!(stats.stamped, 1);
assert_eq!(stats.changed, 1, "the stored verdict was wrong");
let (is_ai, tool, mode, version) = verdict(&db, "stale_trailer");
assert_eq!(is_ai, 1);
assert_eq!(tool.as_deref(), Some("claude"));
assert_eq!(mode, "full_agentic");
assert_eq!(version, DETECTOR_VERSION);
}
#[test]
fn a_generation_1_row_gains_its_detection_method() {
let mut db = Database::open_in_memory().expect("open db");
let footer = "docs: link the website (#5330)\n\n\
🤖🤖🤖 Generated with trusty-mpm — \
https://github.com/bobmatnyc/trusty-tools\n";
db.connection()
.execute(
"INSERT INTO commits \
(sha, author_name, author_email, timestamp, message, repository, \
is_ai_assisted, ai_tool, agentic_mode, ai_detector_version, \
ai_detection_method) \
VALUES ('gen1', 'Ada', 'ada@example.com', '2026-01-01T00:00:00Z', ?1, \
'testrepo', 1, 'trusty-mpm', 'full_agentic', 1, NULL)",
params![footer],
)
.expect("seed a generation-1 row");
let stats = reclassify_stale(&mut db).expect("reclassify");
assert_eq!(stats.stamped, 1);
assert_eq!(
stats.changed, 1,
"the method moved from NULL, so the row changed even though the \
tool and mode did not"
);
let method: Option<String> = db
.connection()
.query_row(
"SELECT ai_detection_method FROM commits WHERE sha = 'gen1'",
[],
|r| r.get(0),
)
.expect("read method");
assert_eq!(method.as_deref(), Some("message"));
let (is_ai, tool, mode, version) = verdict(&db, "gen1");
assert_eq!(
(is_ai, tool.as_deref(), mode.as_str()),
(1, Some("trusty-mpm"), "full_agentic"),
"the rest of the verdict is unchanged"
);
assert_eq!(version, DETECTOR_VERSION);
let again = reclassify_stale(&mut db).expect("reclassify twice");
assert_eq!(again.stamped, 0);
}
#[test]
fn the_scan_uses_the_index_on_a_populated_database() {
let path = std::env::temp_dir().join(format!(
"tga-6748-plan-{}-{:?}.db",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_file(&path);
let db = Database::open(&path).expect("open db");
{
let conn = db.connection();
conn.execute_batch("BEGIN").expect("begin");
let mut insert = conn
.prepare(
"INSERT INTO commits (sha, author_name, author_email, timestamp, \
message, repository, is_ai_assisted, ai_tool, agentic_mode, \
ai_detector_version) \
VALUES (?1, 'Ada', 'ada@example.com', '2026-01-01T00:00:00Z', ?2, \
'testrepo', 0, NULL, 'none', ?3)",
)
.expect("prepare insert");
for i in 0..20_000 {
insert
.execute(params![
format!("sha{i:08}"),
format!("feat: commit {i} with a body long enough to be a real row"),
DETECTOR_VERSION
])
.expect("insert");
}
drop(insert);
conn.execute_batch("COMMIT; ANALYZE;").expect("analyze");
let stat: String = conn
.query_row(
"SELECT stat FROM sqlite_stat1 \
WHERE idx = 'idx_commits_ai_detector_version'",
[],
|r| r.get(0),
)
.expect("the index must have statistics for the plan to be meaningful");
assert!(
stat.starts_with("20000 20000"),
"the pathological statistic this test exists for: {stat}"
);
let mut plan_stmt = conn
.prepare(&format!("EXPLAIN QUERY PLAN {SCAN_SQL}"))
.expect("prepare plan");
let plan = plan_stmt
.query_map(params![DETECTOR_VERSION, RECLASSIFY_BATCH as i64], |r| {
r.get::<_, String>(3)
})
.expect("query plan")
.collect::<std::result::Result<Vec<_>, _>>()
.expect("collect plan")
.join(" | ");
assert!(
plan.contains("USING INDEX idx_commits_ai_detector_version"),
"the scan must be an index range walk: {plan}"
);
assert!(
!plan.contains("SCAN commits"),
"a full table scan reads every message on every collect: {plan}"
);
assert!(
!plan.contains("TEMP B-TREE"),
"ordering by the index's own columns must need no sort: {plan}"
);
}
drop(db);
let _ = std::fs::remove_file(&path);
}
#[test]
fn stale_rows_are_reclassified_and_current_rows_are_not() {
let mut db = Database::open_in_memory().expect("open db");
insert_commit(&db, "stale_a", TRAILER, 0);
insert_commit(&db, "stale_b", "fix: human work\n", 0);
insert_commit(&db, "current", TRAILER, DETECTOR_VERSION);
let calls = Cell::new(0_usize);
let stats = reclassify_stale_with(&mut db, RECLASSIFY_BATCH, |s| {
calls.set(calls.get() + 1);
detect(s)
})
.expect("reclassify");
assert_eq!(calls.get(), 2, "only the two stale rows may be detected");
assert_eq!(stats.stamped, 2);
assert_eq!(stats.changed, 1, "only `stale_a` carries a marker");
assert_eq!(
verdict(&db, "current").0,
0,
"a row at the current generation is left exactly as stored"
);
let again = Cell::new(0_usize);
let second = reclassify_stale_with(&mut db, RECLASSIFY_BATCH, |s| {
again.set(again.get() + 1);
detect(s)
})
.expect("reclassify twice");
assert_eq!(again.get(), 0, "a settled corpus detects nothing");
assert_eq!(second, ReclassifyStats::default());
}
#[test]
fn an_interrupted_pass_resumes_from_where_it_stopped() {
let mut db = Database::open_in_memory().expect("open db");
for i in 0..5 {
insert_commit(&db, &format!("sha{i}"), TRAILER, 0);
}
let mut detector = detect;
let first = reclassify_batch_with(&mut db, 2, &mut detector).expect("one batch");
assert_eq!(first.stamped, 2, "the interrupt lands after one batch");
for sha in ["sha0", "sha1"] {
let (is_ai, tool, mode, version) = verdict(&db, sha);
assert_eq!(
(is_ai, tool.as_deref(), mode.as_str(), version),
(1, Some("claude"), "full_agentic", DETECTOR_VERSION),
"{sha}: verdict and generation are written together or not at all"
);
}
let stale: i64 = db
.connection()
.query_row(
"SELECT COUNT(*) FROM commits WHERE ai_detector_version < ?1",
params![DETECTOR_VERSION],
|r| r.get(0),
)
.expect("count stale");
assert_eq!(stale, 3, "the remainder is still selectable");
let calls = Cell::new(0_usize);
let resumed = reclassify_stale_with(&mut db, 2, |s| {
calls.set(calls.get() + 1);
detect(s)
})
.expect("resume");
assert_eq!(calls.get(), 3, "the resuming pass redoes nothing");
assert_eq!(resumed.stamped, 3);
}
}