use std::collections::HashMap;
use tracing::info;
use crate::collect::collector::CollectionStats;
use crate::collect::linear::{LinearClient, LinearIssue};
use crate::core::config::Config;
use crate::core::db::{Database, WorkItemRow};
const LINEAR_SOURCE: &str = "linear";
const LINEAR_ITEM_TYPE: &str = "Issue";
pub(super) async fn fetch_and_store_linear_issues(
db: &mut Database,
config: &Config,
stats: &mut CollectionStats,
) {
if let Some(linear_cfg) = &config.linear {
if linear_cfg.fetch_on_reference {
match LinearClient::new(linear_cfg) {
Ok(client) => {
let commits: Vec<(String, String)> = {
let conn = db.connection();
let mut stmt = match conn.prepare("SELECT sha, message FROM commits") {
Ok(s) => s,
Err(e) => {
stats.fail_stage(format!("Linear: query commits failed: {e}"));
return;
}
};
let rows = match stmt.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
}) {
Ok(r) => r,
Err(e) => {
stats.fail_stage(format!("Linear: read commits failed: {e}"));
return;
}
};
let mut out = Vec::new();
for r in rows.flatten() {
out.push(r);
}
out
};
let messages: Vec<String> =
commits.iter().map(|(_, msg)| msg.clone()).collect();
let msg_refs: Vec<&str> = messages.iter().map(String::as_str).collect();
let issues = match client
.fetch_referenced_issues(&msg_refs, &linear_cfg.team_keys)
.await
{
Ok(issues) => issues,
Err(e) => {
stats.fail_stage(format!("Linear: fetch issues failed: {e}"));
return;
}
};
for issue in &issues {
info!(
id = %issue.identifier,
state = %issue.state,
team = %issue.team,
"Linear issue fetched"
);
}
match client.store_issues(db, &issues) {
Ok(n) => {
info!(stored = n, "persisted linear_issues rows");
stats.linear_issues_fetched += n;
}
Err(e) => {
stats.fail_stage(format!("Linear: store issues failed: {e}"));
}
}
let commit_refs = build_commit_refs(&commits);
match persist_work_items(db, &issues, &commit_refs) {
Ok(n) => {
info!(
stored = n,
source = LINEAR_SOURCE,
"persisted work_items rows"
);
}
Err(e) => {
stats.fail_stage(format!("Linear: store work_items failed: {e}"));
}
}
}
Err(e) => {
stats.fail_stage(format!("Linear client init failed: {e}"));
}
}
}
}
}
fn build_commit_refs(commits: &[(String, String)]) -> HashMap<String, Vec<String>> {
let mut out = HashMap::new();
for (sha, message) in commits {
let ids = LinearClient::extract_issue_ids(message);
if !ids.is_empty() {
out.insert(sha.clone(), ids);
}
}
out
}
pub fn persist_work_items(
db: &mut Database,
issues: &[LinearIssue],
commit_refs: &HashMap<String, Vec<String>>,
) -> crate::core::Result<usize> {
use crate::core::db::work_items::{link_commit_work_item, upsert_work_item};
use std::collections::HashSet;
if issues.is_empty() {
return Ok(0);
}
let tx = db.connection_mut().transaction()?;
let mut written = 0usize;
for issue in issues {
let row = WorkItemRow {
id: issue.identifier.clone(),
source: LINEAR_SOURCE.to_string(),
title: issue.title.clone(),
status: issue.state.clone(),
item_type: LINEAR_ITEM_TYPE.to_string(),
tags: None,
project: None,
url: Some(issue.url.clone()),
raw_json: serde_json::to_string(issue).ok(),
};
upsert_work_item(&tx, &row)?;
written += 1;
}
let fetched: HashSet<&str> = issues.iter().map(|i| i.identifier.as_str()).collect();
for (sha, ids) in commit_refs {
for id in ids {
if fetched.contains(id.as_str()) {
link_commit_work_item(&tx, sha, id, LINEAR_SOURCE)?;
}
}
}
tx.commit()?;
Ok(written)
}
#[cfg(test)]
mod tests {
use super::*;
fn issue(identifier: &str) -> LinearIssue {
LinearIssue {
identifier: identifier.to_string(),
title: format!("Title for {identifier}"),
state: "In Progress".to_string(),
team: "Engineering".to_string(),
assignee: Some("Alice".to_string()),
priority: 2,
url: format!("https://linear.app/x/issue/{identifier}"),
}
}
fn linear_row_count(db: &Database) -> i64 {
db.connection()
.query_row(
"SELECT COUNT(*) FROM work_items WHERE source = 'linear'",
[],
|r| r.get(0),
)
.expect("count work_items")
}
#[test]
fn persist_work_items_writes_linear_rows() {
let mut db = Database::open_in_memory().expect("db");
let n = persist_work_items(&mut db, &[issue("ENG-1"), issue("FE-42")], &HashMap::new())
.expect("persist");
assert_eq!(n, 2);
assert_eq!(linear_row_count(&db), 2);
let (title, status, item_type, url): (String, String, String, Option<String>) = db
.connection()
.query_row(
"SELECT title, status, item_type, url FROM work_items \
WHERE id = ?1 AND source = 'linear'",
["ENG-1"],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
)
.expect("query row");
assert_eq!(title, "Title for ENG-1");
assert_eq!(status, "In Progress");
assert_eq!(item_type, "Issue");
assert_eq!(url.as_deref(), Some("https://linear.app/x/issue/ENG-1"));
}
#[test]
fn persist_work_items_links_commits() {
let mut db = Database::open_in_memory().expect("db");
let commits = vec![
("sha1".to_string(), "ENG-1: add login".to_string()),
("sha2".to_string(), "chore: no ticket here".to_string()),
];
let refs = build_commit_refs(&commits);
assert_eq!(refs.len(), 1, "only the ticketed commit contributes a ref");
persist_work_items(&mut db, &[issue("ENG-1")], &refs).expect("persist");
let linked: Vec<String> = {
let conn = db.connection();
let mut stmt = conn
.prepare(
"SELECT commit_sha FROM commit_work_items \
WHERE work_item_id = ?1 AND work_item_source = 'linear'",
)
.expect("prepare");
let rows = stmt
.query_map(["ENG-1"], |r| r.get::<_, String>(0))
.expect("query");
rows.flatten().collect()
};
assert_eq!(linked, vec!["sha1".to_string()]);
}
#[test]
fn persist_work_items_skips_unfetched_refs() {
let mut db = Database::open_in_memory().expect("db");
let commits = vec![("sha1".to_string(), "ENG-1 and GONE-9".to_string())];
let refs = build_commit_refs(&commits);
persist_work_items(&mut db, &[issue("ENG-1")], &refs).expect("persist");
let links: i64 = db
.connection()
.query_row("SELECT COUNT(*) FROM commit_work_items", [], |r| r.get(0))
.expect("count links");
assert_eq!(links, 1, "only the fetched identifier is linked");
}
#[test]
fn persist_work_items_is_idempotent() {
let mut db = Database::open_in_memory().expect("db");
let commits = vec![("sha1".to_string(), "ENG-1: work".to_string())];
let refs = build_commit_refs(&commits);
persist_work_items(&mut db, &[issue("ENG-1")], &refs).expect("first");
let mut updated = issue("ENG-1");
updated.state = "Done".to_string();
persist_work_items(&mut db, &[updated], &refs).expect("second");
assert_eq!(linear_row_count(&db), 1);
let status: String = db
.connection()
.query_row(
"SELECT status FROM work_items WHERE id = ?1 AND source = 'linear'",
["ENG-1"],
|r| r.get(0),
)
.expect("query status");
assert_eq!(
status, "Done",
"re-running refreshes rather than duplicates"
);
let links: i64 = db
.connection()
.query_row("SELECT COUNT(*) FROM commit_work_items", [], |r| r.get(0))
.expect("count links");
assert_eq!(links, 1);
}
#[test]
fn persist_work_items_reports_write_failure() {
let mut db = Database::open_in_memory().expect("db");
db.connection()
.execute("DROP TABLE work_items", [])
.expect("drop table");
let err = persist_work_items(&mut db, &[issue("ENG-1")], &HashMap::new())
.expect_err("write against a missing table must fail");
assert!(
err.to_string().contains("work_items"),
"error names the failing table: {err}"
);
}
#[test]
fn persist_work_items_rolls_back_a_partial_write() {
let mut db = Database::open_in_memory().expect("db");
db.connection()
.execute("DROP TABLE commit_work_items", [])
.expect("drop join table");
let commits = vec![("sha1".to_string(), "ENG-1: work".to_string())];
let refs = build_commit_refs(&commits);
let err = persist_work_items(&mut db, &[issue("ENG-1")], &refs)
.expect_err("link against a missing table must fail");
assert!(
err.to_string().contains("commit_work_items"),
"error names the failing table: {err}"
);
assert_eq!(
linear_row_count(&db),
0,
"the issue upsert that already succeeded must roll back with the transaction"
);
}
#[tokio::test]
async fn linear_stage_faults_are_recorded_as_stage_failures() {
use crate::core::config::LinearConfig;
let mut db = Database::open_in_memory().expect("db");
db.connection()
.execute("DROP TABLE commits", [])
.expect("drop commits");
let config = Config {
linear: Some(LinearConfig {
api_key: Some("test-key".to_string()),
fetch_on_reference: true,
..Default::default()
}),
..Default::default()
};
let mut stats = CollectionStats::default();
fetch_and_store_linear_issues(&mut db, &config, &mut stats).await;
assert_eq!(stats.errors.len(), 1, "one fault: {:?}", stats.errors);
assert_eq!(
stats.stage_failures().len(),
1,
"a Linear stage that never wrote must reach the exit code: {:?}",
stats.errors
);
}
#[test]
fn persist_work_items_no_issues_is_a_noop() {
let mut db = Database::open_in_memory().expect("db");
assert_eq!(
persist_work_items(&mut db, &[], &HashMap::new()).expect("persist"),
0
);
assert_eq!(linear_row_count(&db), 0);
}
}