use std::collections::HashSet;
use reqwest::Client;
use rusqlite::params;
use serde::{Deserialize, Serialize};
use trusty_common::credentials::scrub_secrets;
use crate::collect::errors::{CollectError, Result};
use crate::core::config::LinearConfig;
use crate::core::db::Database;
const USER_AGENT_VALUE: &str = "trusty-git-analytics/0.1";
const LINEAR_GRAPHQL_URL: &str = "https://api.linear.app/graphql";
const MAX_ERROR_BODY_CHARS: usize = 500;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LinearIssue {
pub identifier: String,
pub title: String,
pub state: String,
pub team: String,
pub assignee: Option<String>,
pub priority: u8,
pub url: String,
}
pub struct LinearClient {
client: Client,
api_key: String,
endpoint: String,
}
const REDACTED_API_KEY: &str = "<redacted>";
impl std::fmt::Debug for LinearClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LinearClient")
.field("endpoint", &self.endpoint)
.field("api_key", &REDACTED_API_KEY)
.finish_non_exhaustive()
}
}
impl LinearClient {
pub fn new(config: &LinearConfig) -> Result<Self> {
let api_key = resolve_api_key(config)?;
let client = Client::builder()
.user_agent(USER_AGENT_VALUE)
.timeout(std::time::Duration::from_secs(30))
.build()
.map_err(CollectError::Http)?;
Ok(Self {
client,
api_key,
endpoint: LINEAR_GRAPHQL_URL.to_string(),
})
}
#[cfg(test)]
fn with_endpoint(config: &LinearConfig, endpoint: impl Into<String>) -> Result<Self> {
Ok(Self {
endpoint: endpoint.into(),
..Self::new(config)?
})
}
pub async fn fetch_issue(&self, identifier: &str) -> Result<Option<LinearIssue>> {
let query = format!(
r#"query {{
issue(id: "{identifier}") {{
identifier
title
state {{ name }}
team {{ name }}
assignee {{ displayName }}
priority
url
}}
}}"#
);
let body = serde_json::json!({ "query": query });
let resp = self
.client
.post(&self.endpoint)
.header("Authorization", &self.api_key)
.header("Content-Type", "application/json")
.json(&body)
.send()
.await
.map_err(CollectError::Http)?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
return Err(CollectError::LinearApi {
status: status.as_u16(),
identifier: identifier.to_string(),
message: redacted_body_excerpt(&body, &self.api_key),
});
}
let json: serde_json::Value = resp.json().await.map_err(CollectError::Http)?;
if let Some(errors) = json.get("errors") {
if errors.as_array().is_some_and(|a| !a.is_empty()) {
let detail = redacted_body_excerpt(&errors.to_string(), &self.api_key);
tracing::warn!(
identifier = %identifier,
errors = %detail,
"Linear GraphQL errors"
);
return Ok(None);
}
}
let issue_val = &json["data"]["issue"];
if issue_val.is_null() {
return Ok(None);
}
Ok(Some(LinearIssue {
identifier: issue_val["identifier"]
.as_str()
.unwrap_or(identifier)
.to_string(),
title: issue_val["title"].as_str().unwrap_or("").to_string(),
state: issue_val["state"]["name"]
.as_str()
.unwrap_or("Unknown")
.to_string(),
team: issue_val["team"]["name"]
.as_str()
.unwrap_or("Unknown")
.to_string(),
assignee: issue_val["assignee"]["displayName"]
.as_str()
.map(String::from),
priority: issue_val["priority"].as_u64().unwrap_or(0) as u8,
url: issue_val["url"].as_str().unwrap_or("").to_string(),
}))
}
pub fn extract_issue_ids(message: &str) -> Vec<String> {
let re = regex::Regex::new(r"\b([A-Z][A-Z0-9]{0,9}-\d+)\b").expect("valid regex");
let mut seen = HashSet::new();
let mut out = Vec::new();
for cap in re.captures_iter(message) {
let id = cap[1].to_string();
if seen.insert(id.clone()) {
out.push(id);
}
}
out
}
pub async fn fetch_referenced_issues(
&self,
messages: &[&str],
team_filter: &[String],
) -> Result<Vec<LinearIssue>> {
let mut seen = HashSet::new();
let mut all_ids: Vec<String> = Vec::new();
for msg in messages {
for id in Self::extract_issue_ids(msg) {
if seen.insert(id.clone()) {
all_ids.push(id);
}
}
}
let ids: Vec<String> = if team_filter.is_empty() {
all_ids
} else {
all_ids
.into_iter()
.filter(|id| {
let team_key = id.split('-').next().unwrap_or("");
team_filter.iter().any(|t| t.eq_ignore_ascii_case(team_key))
})
.collect()
};
let mut issues = Vec::new();
for id in &ids {
match self.fetch_issue(id).await? {
Some(issue) => issues.push(issue),
None => tracing::debug!("Linear issue not found: {id}"),
}
}
Ok(issues)
}
pub fn store_issues(
&self,
db: &Database,
issues: &[LinearIssue],
) -> crate::core::Result<usize> {
store_linear_issues(db, issues)
}
}
pub fn store_linear_issues(db: &Database, issues: &[LinearIssue]) -> crate::core::Result<usize> {
let conn = db.connection();
let fetched_at = chrono::Utc::now().to_rfc3339();
let mut count = 0usize;
for issue in issues {
let team_key = issue.identifier.split('-').next().unwrap_or("").to_string();
conn.execute(
"INSERT OR REPLACE INTO linear_issues \
(identifier, title, state, team, team_key, assignee, priority, url, fetched_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
params![
issue.identifier,
issue.title,
issue.state,
issue.team,
team_key,
issue.assignee,
issue.priority as i64,
issue.url,
fetched_at,
],
)?;
count += 1;
}
Ok(count)
}
fn redacted_body_excerpt(body: &str, api_key: &str) -> String {
let clean = scrub_secrets(body, &[api_key]);
let trimmed = clean.trim();
match trimmed.char_indices().nth(MAX_ERROR_BODY_CHARS) {
Some((idx, _)) => format!("{}…", &trimmed[..idx]),
None => trimmed.to_string(),
}
}
fn expand_env_var(raw: &str) -> String {
crate::collect::env_expand::expand_env_var(raw)
}
const LINEAR_CREDENTIAL_PROVIDER: &str = "linear";
fn resolve_api_key(config: &LinearConfig) -> Result<String> {
resolve_api_key_with(config, || {
trusty_common::credentials::resolve_key(LINEAR_CREDENTIAL_PROVIDER)
})
}
fn resolve_api_key_with(
config: &LinearConfig,
fallback: impl FnOnce() -> Option<String>,
) -> Result<String> {
let configured = expand_env_var(config.api_key.as_deref().unwrap_or(""));
if !configured.is_empty() {
return Ok(configured);
}
fallback().filter(|k| !k.is_empty()).ok_or_else(|| {
CollectError::Config(
"Linear api_key is required — set `linear.api_key` in the config (a \
`${LINEAR_API_KEY}` reference is expanded), or provide LINEAR_API_KEY \
in the environment, in .env.local, or in the credential store"
.into(),
)
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn extract_issue_ids_finds_linear_patterns() {
let msg = "ENG-123: add login feature, also fixes FE-456";
let ids = LinearClient::extract_issue_ids(msg);
assert!(ids.contains(&"ENG-123".to_string()));
assert!(ids.contains(&"FE-456".to_string()));
}
#[test]
fn extract_issue_ids_deduplicates() {
let msg = "ENG-123 ENG-123 duplicate";
let ids = LinearClient::extract_issue_ids(msg);
assert_eq!(ids.len(), 1);
assert_eq!(ids[0], "ENG-123");
}
#[test]
fn extract_issue_ids_ignores_lowercase_prefix() {
let msg = "abc-123 should not match";
let ids = LinearClient::extract_issue_ids(msg);
assert!(ids.is_empty());
}
#[test]
fn new_rejects_missing_api_key() {
let cfg = LinearConfig::default();
let err = resolve_api_key_with(&cfg, || None).expect_err("should reject empty key");
match err {
CollectError::Config(msg) => assert!(msg.contains("api_key")),
other => panic!("unexpected: {other:?}"),
}
}
#[test]
fn config_api_key_wins_over_the_resolver() {
let cfg = LinearConfig {
api_key: Some("lin_api_from_config".into()),
..LinearConfig::default()
};
let key = resolve_api_key_with(&cfg, || Some("lin_api_from_store".into())).expect("key");
assert_eq!(key, "lin_api_from_config");
}
#[test]
fn an_absent_config_key_falls_back_to_the_resolver() {
let cfg = LinearConfig::default();
let key = resolve_api_key_with(&cfg, || Some("lin_api_from_store".into())).expect("key");
assert_eq!(key, "lin_api_from_store");
let placeholder = LinearConfig {
api_key: Some("${TGA_LINEAR_KEY_THAT_IS_NEVER_SET}".into()),
..LinearConfig::default()
};
let key =
resolve_api_key_with(&placeholder, || Some("lin_api_from_store".into())).expect("key");
assert_eq!(key, "lin_api_from_store");
}
#[test]
fn an_empty_resolver_answer_is_not_a_key() {
let cfg = LinearConfig::default();
let err = resolve_api_key_with(&cfg, || Some(String::new()))
.expect_err("empty is not a credential");
match err {
CollectError::Config(msg) => assert!(msg.contains("LINEAR_API_KEY"), "{msg}"),
other => panic!("unexpected: {other:?}"),
}
}
fn sample_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}"),
}
}
#[test]
fn store_linear_issues_inserts_rows() {
let db = Database::open_in_memory().expect("db");
let issues = vec![sample_issue("ENG-1"), sample_issue("FE-42")];
let n = store_linear_issues(&db, &issues).expect("store");
assert_eq!(n, 2);
let conn = db.connection();
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM linear_issues", [], |r| r.get(0))
.expect("count");
assert_eq!(count, 2);
let (identifier, team_key, priority): (String, String, i64) = conn
.query_row(
"SELECT identifier, team_key, priority FROM linear_issues WHERE identifier = ?1",
["ENG-1"],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("query");
assert_eq!(identifier, "ENG-1");
assert_eq!(team_key, "ENG");
assert_eq!(priority, 2);
}
#[test]
fn store_linear_issues_is_idempotent_on_identifier() {
let db = Database::open_in_memory().expect("db");
let mut issue = sample_issue("ENG-9");
store_linear_issues(&db, &[issue.clone()]).expect("first");
issue.state = "Done".to_string();
issue.assignee = Some("Bob".to_string());
store_linear_issues(&db, &[issue]).expect("second");
let conn = db.connection();
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM linear_issues", [], |r| r.get(0))
.expect("count");
assert_eq!(count, 1);
let (state, assignee): (String, Option<String>) = conn
.query_row(
"SELECT state, assignee FROM linear_issues WHERE identifier = ?1",
["ENG-9"],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.expect("query");
assert_eq!(state, "Done");
assert_eq!(assignee.as_deref(), Some("Bob"));
}
#[test]
fn store_linear_issues_handles_missing_assignee() {
let db = Database::open_in_memory().expect("db");
let mut issue = sample_issue("OPS-7");
issue.assignee = None;
store_linear_issues(&db, &[issue]).expect("store");
let conn = db.connection();
let assignee: Option<String> = conn
.query_row(
"SELECT assignee FROM linear_issues WHERE identifier = ?1",
["OPS-7"],
|r| r.get(0),
)
.expect("query");
assert!(assignee.is_none());
}
#[test]
fn migration_v2_creates_linear_issues_table() {
let db = Database::open_in_memory().expect("db");
let conn = db.connection();
let name: String = conn
.query_row(
"SELECT name FROM sqlite_master WHERE type='table' AND name='linear_issues'",
[],
|r| r.get(0),
)
.expect("table exists");
assert_eq!(name, "linear_issues");
assert!(db.schema_version().expect("version") >= 2);
}
const FAKE_API_KEY: &str = "lin_api_averyrealisticlookingkey0123456789";
#[test]
fn redacted_body_excerpt_keeps_short_input() {
assert_eq!(
redacted_body_excerpt(" {\"errors\":[]} ", FAKE_API_KEY),
"{\"errors\":[]}"
);
}
#[test]
fn redacted_body_excerpt_clips_long_input() {
let out = redacted_body_excerpt(&"x".repeat(MAX_ERROR_BODY_CHARS + 50), FAKE_API_KEY);
assert_eq!(out.chars().count(), MAX_ERROR_BODY_CHARS + 1);
assert!(out.ends_with('…'));
}
#[test]
fn redacted_body_excerpt_scrubs_before_truncating() {
let pad = "x".repeat(MAX_ERROR_BODY_CHARS - 30);
let body = format!("{pad}{FAKE_API_KEY} trailing detail");
let out = redacted_body_excerpt(&body, FAKE_API_KEY);
assert!(!out.contains(FAKE_API_KEY), "whole key survived: {out}");
assert!(
!out.contains(&FAKE_API_KEY[..30]),
"a prefix of the key survived the cut — truncation ran first: {out}"
);
assert!(out.contains("[REDACTED]"), "key was not scrubbed: {out}");
}
const PROBE_ENDPOINT: &str = "http://endpoint.invalid/graphql";
fn client_holding(key: &str) -> LinearClient {
LinearClient {
client: Client::new(),
api_key: key.to_string(),
endpoint: PROBE_ENDPOINT.to_string(),
}
}
#[test]
fn debug_never_renders_the_api_key() {
let cases: &[(&str, &str)] = &[
(FAKE_API_KEY, "the lin_-prefixed key production expects"),
(
"9f3Kq7Zt2Wm4Bx8Lv6Nc1Rd5Ph0Sj",
"no prefix: entropy up front",
),
("ab7Q", "exactly a four-character head"),
("x9", "shorter than a head"),
("", "empty — unreachable via new(), guarded anyway"),
];
for (key, why) in cases {
let client = client_holding(key);
let compact = format!("{client:?}");
let pretty = format!("{client:#?}");
for rendered in [&compact, &pretty] {
if !key.is_empty() {
assert!(
!rendered.contains(key),
"{why}: the whole key reached Debug output: {rendered}"
);
let head: String = key.chars().take(4).collect();
assert!(
!rendered.contains(&head),
"{why}: a leading fragment of the key survived: {rendered}"
);
}
assert!(
rendered.contains(REDACTED_API_KEY),
"{why}: the key field was not masked: {rendered}"
);
assert!(
rendered.contains("endpoint.invalid"),
"{why}: redaction must not cost the endpoint, the field \
worth debugging: {rendered}"
);
}
}
}
#[tokio::test]
#[tracing_test::traced_test]
async fn graphql_errors_are_scrubbed_before_they_reach_the_log() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"errors": [{
"message": format!("API key {FAKE_API_KEY} lacks the read scope")
}]
})))
.mount(&server)
.await;
let got = mock_client(&server.uri())
.fetch_issue("ENG-1")
.await
.expect("a 200 carrying GraphQL errors is still a successful call");
assert!(got.is_none(), "the #5665 control flow is deliberately kept");
assert!(
!logs_contain(FAKE_API_KEY),
"the key reached the operator's terminal"
);
assert!(logs_contain("[REDACTED]"), "the key was not scrubbed");
assert!(
logs_contain("lacks the read scope"),
"redaction must not cost the reader Linear's diagnosis"
);
}
fn mock_client(endpoint: &str) -> LinearClient {
let cfg = LinearConfig {
api_key: Some(FAKE_API_KEY.into()),
..Default::default()
};
LinearClient::with_endpoint(&cfg, endpoint).expect("client builds")
}
const AUTH_ERROR_BODY: &str = r#"{"errors":[{"message":"Authentication required, not authenticated","extensions":{"type":"authentication error","code":"AUTHENTICATION_ERROR","statusCode":401,"userPresentableMessage":"You need to authenticate to access this operation."}}]}"#;
#[tokio::test]
async fn fetch_issue_errors_on_auth_failure() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(
ResponseTemplate::new(401).set_body_raw(AUTH_ERROR_BODY, "application/json"),
)
.mount(&server)
.await;
let err = mock_client(&server.uri())
.fetch_issue("ENG-1")
.await
.expect_err("a 401 must not read as an absent issue");
match err {
CollectError::LinearApi {
status,
identifier,
message,
} => {
assert_eq!(status, 401);
assert_eq!(identifier, "ENG-1");
assert!(
message.contains("You need to authenticate"),
"Linear's own diagnosis must survive into the error: {message}"
);
}
other => panic!("expected LinearApi, got {other:?}"),
}
}
#[tokio::test]
async fn an_api_key_echoed_in_the_error_body_never_reaches_the_message() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
let echoing_body = format!(
r#"{{"errors":[{{"message":"API key {FAKE_API_KEY} is not valid for this workspace"}}]}}"#
);
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(401).set_body_raw(echoing_body, "application/json"))
.mount(&server)
.await;
let err = mock_client(&server.uri())
.fetch_issue("ENG-1")
.await
.expect_err("a 401 must surface");
let rendered = err.to_string();
assert!(
!rendered.contains(FAKE_API_KEY),
"the key reached the error message: {rendered}"
);
assert!(
rendered.contains("[REDACTED]"),
"the key was not scrubbed: {rendered}"
);
assert!(
rendered.contains("is not valid for this workspace"),
"redaction must not cost the reader Linear's diagnosis: {rendered}"
);
}
#[tokio::test]
async fn fetch_issue_errors_on_server_failure() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(500).set_body_raw("upstream down", "text/plain"))
.mount(&server)
.await;
let err = mock_client(&server.uri())
.fetch_issue("ENG-1")
.await
.expect_err("a 500 must surface");
assert!(
matches!(err, CollectError::LinearApi { status: 500, .. }),
"expected a 500 LinearApi, got {err:?}"
);
}
#[tokio::test]
async fn fetch_issue_returns_none_for_absent_issue() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"data": { "issue": null }
})))
.mount(&server)
.await;
let got = mock_client(&server.uri())
.fetch_issue("ENG-404")
.await
.expect("an absent issue is a successful answer");
assert!(got.is_none());
}
#[tokio::test]
async fn fetch_referenced_issues_propagates_auth_failure() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(
ResponseTemplate::new(401).set_body_raw(AUTH_ERROR_BODY, "application/json"),
)
.mount(&server)
.await;
let err = mock_client(&server.uri())
.fetch_referenced_issues(&["ENG-1: work", "FE-2: more"], &[])
.await
.expect_err("an invalid key must reach the caller");
assert!(
matches!(err, CollectError::LinearApi { status: 401, .. }),
"expected a 401 LinearApi, got {err:?}"
);
assert_eq!(
server.received_requests().await.map(|r| r.len()),
Some(1),
"the walk stops at the first failure instead of retrying every id"
);
}
#[tokio::test]
async fn fetch_referenced_issues_skips_absent_issues() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"data": { "issue": null }
})))
.mount(&server)
.await;
let issues = mock_client(&server.uri())
.fetch_referenced_issues(&["ENG-1 and FE-2"], &[])
.await
.expect("absent issues are not a failure");
assert!(issues.is_empty());
}
#[tokio::test]
async fn fetch_issue_live() {
let key = match std::env::var("LINEAR_API_KEY") {
Ok(k) => k,
Err(_) => {
eprintln!("SKIP: set LINEAR_API_KEY to run");
return;
}
};
let config = LinearConfig {
api_key: Some(key),
..Default::default()
};
let client = LinearClient::new(&config).expect("client");
let result = client.fetch_issue("ENG-1").await;
assert!(
result.is_ok(),
"fetch must not error — a revoked key lands here: {result:?}"
);
println!("Result: {result:?}");
}
}