use std::collections::HashSet;
use chrono::{DateTime, Utc};
use reqwest::Client;
use rusqlite::params;
use serde::{Deserialize, Serialize};
use trusty_common::credentials::scrub_secrets;
use crate::collect::errors::{CollectError, Result};
use crate::collect::jira::retry::{with_retry, RetryBudget, RetryPolicy};
use crate::collect::linear::sync;
use crate::collect::ticket::is_non_ticket_identifier;
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 LINEAR_PAGE_BUDGET_SLACK: usize = 8;
const MAX_ERROR_BODY_CHARS: usize = 500;
#[derive(Debug, Clone, PartialEq, 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,
#[serde(default)]
pub created_at: Option<DateTime<Utc>>,
#[serde(default)]
pub updated_at: Option<DateTime<Utc>>,
#[serde(default)]
pub started_at: Option<DateTime<Utc>>,
#[serde(default)]
pub completed_at: Option<DateTime<Utc>>,
#[serde(default)]
pub canceled_at: Option<DateTime<Utc>>,
}
pub struct LinearClient {
client: Client,
api_key: String,
endpoint: String,
retry: RetryPolicy,
budget: RetryBudget,
}
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)?;
let retry = RetryPolicy::default();
let budget = RetryBudget::new(&retry);
Ok(Self {
client,
api_key,
endpoint: LINEAR_GRAPHQL_URL.to_string(),
retry,
budget,
})
}
#[must_use]
pub fn with_retry_policy(mut self, policy: RetryPolicy) -> Self {
self.budget = RetryBudget::new(&policy);
self.retry = policy;
self
}
#[cfg(test)]
pub(crate) 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
createdAt
updatedAt
startedAt
completedAt
canceledAt
}}
}}"#
);
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(parse_issue_node(identifier, issue_val)))
}
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 is_non_ticket_identifier(&id) {
continue;
}
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 async fn fetch_team_issues_page(
&self,
team_key: &str,
since: Option<DateTime<Utc>>,
after: Option<&str>,
page_size: usize,
page_number: usize,
) -> Result<LinearIssuesPage> {
with_retry("linear issues page", &self.retry, &self.budget, || {
self.send_team_issues_page(team_key, since, after, page_size, page_number)
})
.await
}
async fn send_team_issues_page(
&self,
team_key: &str,
since: Option<DateTime<Utc>>,
after: Option<&str>,
page_size: usize,
page_number: usize,
) -> Result<LinearIssuesPage> {
const QUERY: &str = r#"query($first: Int!, $after: String, $filter: IssueFilter, $orderBy: PaginationOrderBy) {
issues(first: $first, after: $after, filter: $filter, orderBy: $orderBy) {
nodes {
identifier
title
state { name }
team { name key }
assignee { displayName }
priority
url
createdAt
updatedAt
startedAt
completedAt
canceledAt
}
pageInfo { hasNextPage endCursor }
}
}"#;
let variables = serde_json::json!({
"first": page_size,
"after": after,
"filter": sync::build_issues_filter(team_key, since),
"orderBy": "updatedAt",
});
let body = serde_json::json!({ "query": QUERY, "variables": variables });
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 == reqwest::StatusCode::TOO_MANY_REQUESTS
|| status == reqwest::StatusCode::SERVICE_UNAVAILABLE
{
let retry_after = resp
.headers()
.get(reqwest::header::RETRY_AFTER)
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse::<u64>().ok())
.map(std::time::Duration::from_secs);
return Err(CollectError::Throttled {
status: status.as_u16(),
retry_after,
});
}
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
return Err(CollectError::LinearBulkApi {
status: status.as_u16(),
team_key: team_key.to_string(),
page: page_number,
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);
return Err(CollectError::LinearBulkApi {
status: status.as_u16(),
team_key: team_key.to_string(),
page: page_number,
message: detail,
});
}
}
let nodes = json["data"]["issues"]["nodes"]
.as_array()
.cloned()
.unwrap_or_default();
let issues = nodes
.iter()
.map(|node| parse_issue_node(node["identifier"].as_str().unwrap_or(""), node))
.collect();
let has_next_page = json["data"]["issues"]["pageInfo"]["hasNextPage"]
.as_bool()
.unwrap_or(false);
let end_cursor = json["data"]["issues"]["pageInfo"]["endCursor"]
.as_str()
.map(String::from);
Ok(LinearIssuesPage {
issues,
has_next_page,
end_cursor,
})
}
pub async fn fetch_team_issues(
&self,
team_key: &str,
since: Option<DateTime<Utc>>,
max_issues: usize,
) -> Result<(Vec<LinearIssue>, bool)> {
const PAGE_SIZE: usize = 50;
let max_pages = max_issues.div_ceil(PAGE_SIZE) + LINEAR_PAGE_BUDGET_SLACK;
let mut issues = Vec::new();
let mut after: Option<String> = None;
let mut page_number = 0usize;
loop {
page_number += 1;
let page = self
.fetch_team_issues_page(team_key, since, after.as_deref(), PAGE_SIZE, page_number)
.await?;
issues.extend(page.issues);
if issues.len() >= max_issues {
issues.truncate(max_issues);
return Ok((issues, true));
}
if !page.has_next_page || page.end_cursor.is_none() {
return Ok((issues, false));
}
if page_number >= max_pages {
return Err(CollectError::PagingBudgetExceeded {
endpoint: "linear/issues",
key: team_key.to_string(),
pages: page_number,
});
}
after = page.end_cursor;
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct LinearIssuesPage {
pub issues: Vec<LinearIssue>,
pub has_next_page: bool,
pub end_cursor: Option<String>,
}
fn parse_issue_node(identifier_fallback: &str, node: &serde_json::Value) -> LinearIssue {
let parse_dt = |field: &str| -> Option<DateTime<Utc>> {
node[field]
.as_str()
.and_then(|s| DateTime::parse_from_rfc3339(s).ok())
.map(|d| d.with_timezone(&Utc))
};
LinearIssue {
identifier: node["identifier"]
.as_str()
.unwrap_or(identifier_fallback)
.to_string(),
title: node["title"].as_str().unwrap_or("").to_string(),
state: node["state"]["name"]
.as_str()
.unwrap_or("Unknown")
.to_string(),
team: node["team"]["name"]
.as_str()
.unwrap_or("Unknown")
.to_string(),
assignee: node["assignee"]["displayName"].as_str().map(String::from),
priority: node["priority"].as_u64().unwrap_or(0) as u8,
url: node["url"].as_str().unwrap_or("").to_string(),
created_at: parse_dt("createdAt"),
updated_at: parse_dt("updatedAt"),
started_at: parse_dt("startedAt"),
completed_at: parse_dt("completedAt"),
canceled_at: parse_dt("canceledAt"),
}
}
pub fn store_linear_issues(db: &Database, issues: &[LinearIssue]) -> crate::core::Result<usize> {
let conn = db.connection();
let tx = conn.unchecked_transaction()?;
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();
tx.execute(
"INSERT OR REPLACE INTO linear_issues \
(identifier, title, state, team, team_key, assignee, priority, url, fetched_at, \
created_at, updated_at, started_at, completed_at, canceled_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)",
params![
issue.identifier,
issue.title,
issue.state,
issue.team,
team_key,
issue.assignee,
issue.priority as i64,
issue.url,
fetched_at,
issue.created_at.map(|d| d.to_rfc3339()),
issue.updated_at.map(|d| d.to_rfc3339()),
issue.started_at.map(|d| d.to_rfc3339()),
issue.completed_at.map(|d| d.to_rfc3339()),
issue.canceled_at.map(|d| d.to_rfc3339()),
],
)?;
count += 1;
}
tx.commit()?;
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 extract_issue_ids_rejects_non_ticket_identifiers() {
for msg in [
"docs: DOC-67 §5 sweep order",
"spec: DOC-70 board axis",
"docs: amend ADR-0029 point 3",
"docs: supersede ADR-0038",
"fix: a multi-byte UTF-8 name breaks the stem",
"chore: verify the artifact against the published SHA-256",
"feat: strip ECMA-48 control sequences",
"chore: avoid reintroducing RUSTSEC-2026-0187",
"fix: ISO-8601 offsets were dropped",
"docs: RFC-2119 keyword sweep",
"chore: drop the MD-5 fallback",
"chore: patch CVE-2024-3094",
"docs: map the finding to CWE-79",
"build: mirror the AL2023/GCC-11 desync",
"chore: rotate to RSA-4096 and AES-256",
"docs: SPEC-14 tightened",
] {
assert_eq!(
LinearClient::extract_issue_ids(msg),
Vec::<String>::new(),
"message: {msg:?}"
);
}
}
#[test]
fn extract_issue_ids_keeps_genuine_tracker_ids() {
for (msg, want) in [
("ABC-123: add the thing", "ABC-123"),
("fix: land GH-12", "GH-12"),
("PROJ-9 tighten the check", "PROJ-9"),
("ENG-123: add login feature", "ENG-123"),
("also fixes FE-456", "FE-456"),
("WI-1 spike", "WI-1"),
("AC-1 acceptance", "AC-1"),
("CREDPANEL-01 wiring", "CREDPANEL-01"),
] {
assert_eq!(
LinearClient::extract_issue_ids(msg),
vec![want.to_string()],
"message: {msg:?}"
);
}
assert!(LinearClient::extract_issue_ids("closes #1234").is_empty());
assert_eq!(
crate::collect::ticket::extract_ticket_id("closes #1234"),
Some("#1234".to_string())
);
}
#[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}"),
created_at: None,
updated_at: None,
started_at: None,
completed_at: None,
canceled_at: None,
}
}
#[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 {
let retry = RetryPolicy::default();
let budget = RetryBudget::new(&retry);
LinearClient {
client: Client::new(),
api_key: key.to_string(),
endpoint: PROBE_ENDPOINT.to_string(),
retry,
budget,
}
}
#[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:?}");
}
mod bulk_sync {
use super::*;
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
fn node(identifier: &str, updated_at: &str) -> serde_json::Value {
serde_json::json!({
"identifier": identifier,
"title": format!("Title for {identifier}"),
"state": {"name": "In Progress"},
"team": {"name": "Engineering", "key": "ENG"},
"assignee": {"displayName": "Alice"},
"priority": 2,
"url": format!("https://linear.app/x/issue/{identifier}"),
"createdAt": "2026-01-01T00:00:00.000Z",
"updatedAt": updated_at,
"startedAt": "2026-01-02T00:00:00.000Z",
"completedAt": serde_json::Value::Null,
"canceledAt": serde_json::Value::Null,
})
}
fn page_response(
nodes: Vec<serde_json::Value>,
has_next: bool,
cursor: Option<&str>,
) -> ResponseTemplate {
ResponseTemplate::new(200).set_body_json(serde_json::json!({
"data": {
"issues": {
"nodes": nodes,
"pageInfo": {"hasNextPage": has_next, "endCursor": cursor},
}
}
}))
}
#[tokio::test]
async fn fetch_team_issues_walks_every_page() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(wiremock::matchers::body_string_contains("\"after\":null"))
.respond_with(page_response(
vec![node("ENG-1", "2026-01-01T00:01:00.000Z")],
true,
Some("cursor-1"),
))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(wiremock::matchers::body_string_contains("cursor-1"))
.respond_with(page_response(
vec![node("ENG-2", "2026-01-01T00:02:00.000Z")],
false,
None,
))
.mount(&server)
.await;
let client = mock_client(&server.uri());
let (issues, truncated) = client
.fetch_team_issues("ENG", None, 10_000)
.await
.expect("walk succeeds");
let ids: Vec<&str> = issues.iter().map(|i| i.identifier.as_str()).collect();
assert_eq!(ids, vec!["ENG-1", "ENG-2"]);
assert!(!truncated);
}
#[tokio::test]
async fn fetch_team_issues_stops_at_max_issues() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(page_response(
vec![
node("ENG-1", "2026-01-01T00:01:00.000Z"),
node("ENG-2", "2026-01-01T00:02:00.000Z"),
node("ENG-3", "2026-01-01T00:03:00.000Z"),
],
true,
Some("cursor-1"),
))
.mount(&server)
.await;
let client = mock_client(&server.uri());
let (issues, truncated) = client
.fetch_team_issues("ENG", None, 2)
.await
.expect("walk succeeds");
assert_eq!(issues.len(), 2);
assert!(truncated);
}
#[tokio::test]
async fn fetch_team_issues_page_maps_timestamps() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(page_response(
vec![node("ENG-1", "2026-01-01T00:01:00.000Z")],
false,
None,
))
.mount(&server)
.await;
let page = mock_client(&server.uri())
.fetch_team_issues_page("ENG", None, None, 50, 1)
.await
.expect("page fetch succeeds");
let issue = &page.issues[0];
assert_eq!(
issue.created_at,
Some(
DateTime::parse_from_rfc3339("2026-01-01T00:00:00.000Z")
.unwrap()
.with_timezone(&Utc)
)
);
assert_eq!(
issue.updated_at,
Some(
DateTime::parse_from_rfc3339("2026-01-01T00:01:00.000Z")
.unwrap()
.with_timezone(&Utc)
)
);
assert!(issue.started_at.is_some());
assert_eq!(issue.completed_at, None);
assert_eq!(issue.canceled_at, None);
}
#[tokio::test]
async fn fetch_team_issues_page_handles_an_empty_team() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(page_response(vec![], false, None))
.mount(&server)
.await;
let page = mock_client(&server.uri())
.fetch_team_issues_page("ENG", None, None, 50, 1)
.await
.expect("empty page is not an error");
assert!(page.issues.is_empty());
assert!(!page.has_next_page);
}
#[tokio::test]
async fn fetch_team_issues_page_errors_on_non_2xx() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(401).set_body_raw(
r#"{"errors":[{"message":"Authentication required"}]}"#,
"application/json",
))
.mount(&server)
.await;
let err = mock_client(&server.uri())
.fetch_team_issues_page("ENG", None, None, 50, 1)
.await
.expect_err("a 401 must not read as an empty team");
match err {
CollectError::LinearBulkApi {
status,
team_key,
page,
..
} => {
assert_eq!(status, 401);
assert_eq!(team_key, "ENG");
assert_eq!(page, 1);
}
other => panic!("expected LinearBulkApi, got {other:?}"),
}
}
fn fast_policy() -> RetryPolicy {
RetryPolicy {
max_attempts: 3,
base_delay: std::time::Duration::from_millis(1),
max_delay: std::time::Duration::from_millis(1),
max_total_delay: std::time::Duration::from_millis(100),
}
}
#[tokio::test]
async fn fetch_team_issues_page_retries_a_429_then_succeeds() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use wiremock::{Request, Respond};
struct OnceThrottled {
calls: Arc<AtomicUsize>,
}
impl Respond for OnceThrottled {
fn respond(&self, _request: &Request) -> ResponseTemplate {
if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
ResponseTemplate::new(429)
.insert_header("Retry-After", "0")
.set_body_raw(
r#"{"errors":[{"message":"rate limited"}]}"#,
"application/json",
)
} else {
page_response(vec![node("ENG-1", "2026-01-01T00:01:00.000Z")], false, None)
}
}
}
let server = MockServer::start().await;
let calls = Arc::new(AtomicUsize::new(0));
Mock::given(method("POST"))
.respond_with(OnceThrottled {
calls: Arc::clone(&calls),
})
.mount(&server)
.await;
let client = mock_client(&server.uri()).with_retry_policy(fast_policy());
let page = client
.fetch_team_issues_page("ENG", None, None, 50, 1)
.await
.expect("the 429 is retried, not surfaced");
assert_eq!(page.issues.len(), 1);
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"expected exactly one retry (429 then 200)"
);
}
#[tokio::test]
async fn fetch_team_issues_errors_when_the_page_budget_is_exhausted() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(page_response(vec![], true, Some("always-more")))
.mount(&server)
.await;
let client = mock_client(&server.uri());
let err = client
.fetch_team_issues("ENG", None, 10)
.await
.expect_err("an endless hasNextPage:true must not loop forever");
match err {
CollectError::PagingBudgetExceeded {
endpoint,
key,
pages,
} => {
assert_eq!(endpoint, "linear/issues");
assert_eq!(key, "ENG");
assert!(pages > 0);
}
other => panic!("expected PagingBudgetExceeded, got {other:?}"),
}
}
}
}