use std::sync::Mutex;
use reqwest::header::{HeaderMap, HeaderValue, ACCEPT, USER_AGENT};
use serde::Deserialize;
use serde_json::{json, Value};
use tracing::debug;
use chrono_tz::Tz;
use crate::collect::errors::{CollectError, Result};
use crate::collect::jira::http::{expand_credential, get_json, post_json, Credentials};
use crate::collect::jira::jql_time::parse_timezone;
use crate::collect::jira::model::{
ChangelogIssue, ChangelogSearchResponse, ChangelogWalk, CommentSearchResponse, JiraComment,
};
use crate::collect::jira::paging::{KeysetPager, PagedItem};
use crate::collect::jira::retry::{with_retry, RetryBudget, RetryPolicy};
use crate::collect::jira::sync::{build_jql, SyncScope};
use crate::core::config::JiraConfig;
#[path = "changelog.rs"]
mod changelog;
const COMMENT_PAGE_SIZE: usize = 100;
const MAX_COMMENT_PAGES: usize = 200;
const USER_AGENT_VALUE: &str = "trusty-git-analytics/0.1";
const SEARCH_PAGE_SIZE: usize = 50;
pub struct JiraClient {
client: reqwest::Client,
base_url: String,
credentials: Option<Credentials>,
project_key: String,
story_point_field: Mutex<Option<Option<String>>>,
configured_timezone: Option<String>,
account_timezone: Mutex<Option<Tz>>,
retry: RetryPolicy,
budget: RetryBudget,
}
#[derive(Debug, Clone)]
pub struct JiraIssue {
pub key: String,
pub summary: String,
pub status: String,
pub issue_type: String,
pub story_points: Option<f64>,
}
#[derive(Debug, Deserialize)]
struct ApiIssue {
key: String,
fields: ApiFields,
}
#[derive(Debug, Deserialize)]
struct ApiFields {
#[serde(default)]
summary: String,
status: ApiNamed,
#[serde(rename = "issuetype")]
issue_type: ApiNamed,
#[serde(flatten)]
extra: std::collections::HashMap<String, Value>,
}
#[derive(Debug, Deserialize)]
struct ApiNamed {
name: String,
}
#[derive(Debug, Deserialize)]
struct FieldDescriptor {
id: String,
name: String,
}
#[derive(Debug, Deserialize)]
struct MyselfResponse {
#[serde(rename = "timeZone", default)]
time_zone: Option<String>,
}
#[derive(Debug, Deserialize)]
struct SearchResponse {
issues: Vec<ApiIssue>,
#[serde(default)]
total: u64,
}
impl JiraClient {
pub fn new(config: &JiraConfig) -> Result<Self> {
let base = config
.url
.as_ref()
.ok_or_else(|| CollectError::Config("jira.url is required".into()))?
.trim_end_matches('/')
.to_string();
let mut headers = HeaderMap::new();
headers.insert(USER_AGENT, HeaderValue::from_static(USER_AGENT_VALUE));
headers.insert(ACCEPT, HeaderValue::from_static("application/json"));
let client = reqwest::Client::builder()
.default_headers(headers)
.timeout(std::time::Duration::from_secs(30))
.build()?;
let credentials = match (&config.username, &config.token) {
(Some(u), Some(t)) => Some((
expand_credential("jira.username", u)?,
expand_credential("jira.token", t)?,
)),
_ => None,
};
let retry = RetryPolicy::default();
Ok(Self {
client,
base_url: base,
credentials,
project_key: config.project_key.clone().unwrap_or_default(),
story_point_field: Mutex::new(None),
configured_timezone: config.timezone.clone(),
account_timezone: Mutex::new(None),
budget: RetryBudget::new(&retry),
retry,
})
}
#[must_use]
pub fn with_retry_policy(mut self, policy: RetryPolicy) -> Self {
self.budget = RetryBudget::new(&policy);
self.retry = policy;
self
}
pub async fn account_timezone(&self) -> Result<Tz> {
{
let guard = self
.account_timezone
.lock()
.map_err(|e| CollectError::Config(format!("timezone cache poisoned: {e}")))?;
if let Some(tz) = *guard {
return Ok(tz);
}
}
let tz = match &self.configured_timezone {
Some(name) => parse_timezone(name)?,
None => {
let url = format!("{}/rest/api/3/myself", self.base_url);
debug!(url = %url, "GET (account timezone)");
let me: MyselfResponse =
with_retry("myself", &self.retry, &self.budget, || self.get(&url))
.await
.map_err(|e| {
CollectError::Config(format!(
"could not determine the JIRA account timezone from \
GET /rest/api/3/myself ({e}). JQL date literals are \
evaluated in the account's timezone, so tga refuses to \
guess — set `jira.timezone` in config.yaml (e.g. `UTC`)."
))
})?;
let name = me.time_zone.ok_or_else(|| {
CollectError::Config(
"the JIRA account reports no `timeZone`; set `jira.timezone` in \
config.yaml so JQL date bounds can be rendered correctly."
.to_string(),
)
})?;
parse_timezone(&name)?
}
};
let mut guard = self
.account_timezone
.lock()
.map_err(|e| CollectError::Config(format!("timezone cache poisoned: {e}")))?;
*guard = Some(tz);
Ok(tz)
}
pub async fn fetch_issue(&self, key: &str) -> Result<Option<JiraIssue>> {
let url = format!("{}/rest/api/3/issue/{}", self.base_url, key);
debug!(url = %url, "GET");
let mut req = self.client.get(&url);
if let Some((user, token)) = &self.credentials {
req = req.basic_auth(user, Some(token));
}
let resp = req.send().await?;
if resp.status() == reqwest::StatusCode::NOT_FOUND {
return Ok(None);
}
let resp = resp.error_for_status()?;
let issue: ApiIssue = resp.json().await?;
let story_field = self.get_story_point_field().await?;
Ok(Some(Self::convert_issue(issue, story_field.as_deref())))
}
pub fn project_key(&self) -> &str {
&self.project_key
}
pub async fn search_issues(&self, jql: &str, max_results: usize) -> Result<Vec<JiraIssue>> {
let url = format!("{}/rest/api/3/search", self.base_url);
let story_field = self.get_story_point_field().await?;
let fields: Vec<String> = match &story_field {
Some(key) => vec![
"summary".into(),
"status".into(),
"issuetype".into(),
key.clone(),
],
None => vec!["*all".into()],
};
let mut out: Vec<JiraIssue> = Vec::new();
let mut start_at = 0u64;
loop {
let remaining = max_results.saturating_sub(out.len());
if remaining == 0 {
break;
}
let page_size = remaining.min(SEARCH_PAGE_SIZE);
let body = json!({
"jql": jql,
"startAt": start_at,
"maxResults": page_size,
"fields": fields,
});
debug!(url = %url, %jql, start_at, "POST");
let mut req = self.client.post(&url).json(&body);
if let Some((user, token)) = &self.credentials {
req = req.basic_auth(user, Some(token));
}
let resp = req.send().await?.error_for_status()?;
let parsed: SearchResponse = resp.json().await?;
let n = parsed.issues.len();
for issue in parsed.issues {
out.push(Self::convert_issue(issue, story_field.as_deref()));
if out.len() >= max_results {
break;
}
}
if n < page_size {
break;
}
start_at += n as u64;
if start_at >= parsed.total {
break;
}
}
Ok(out)
}
pub async fn get_story_point_field(&self) -> Result<Option<String>> {
{
let guard = self
.story_point_field
.lock()
.map_err(|e| CollectError::Config(format!("story-point cache poisoned: {e}")))?;
if let Some(cached) = guard.as_ref() {
return Ok(cached.clone());
}
}
let url = format!("{}/rest/api/3/field", self.base_url);
debug!(url = %url, "GET");
let mut req = self.client.get(&url);
if let Some((user, token)) = &self.credentials {
req = req.basic_auth(user, Some(token));
}
let resp = req.send().await?.error_for_status()?;
let fields: Vec<FieldDescriptor> = resp.json().await?;
let found = fields
.into_iter()
.find(|f| {
let n = f.name.to_ascii_lowercase();
n == "story points" || n == "story point estimate"
})
.map(|f| f.id);
let mut guard = self
.story_point_field
.lock()
.map_err(|e| CollectError::Config(format!("story-point cache poisoned: {e}")))?;
*guard = Some(found.clone());
Ok(found)
}
fn convert_issue(api: ApiIssue, story_field_key: Option<&str>) -> JiraIssue {
let story_points =
story_field_key.and_then(|key| api.fields.extra.get(key).and_then(|v| v.as_f64()));
JiraIssue {
key: api.key,
summary: api.fields.summary,
status: api.fields.status.name,
issue_type: api.fields.issue_type.name,
story_points,
}
}
pub async fn search_with_changelog(
&self,
scope: &SyncScope,
max_results: usize,
) -> Result<ChangelogWalk> {
let url = format!("{}/rest/api/3/search", self.base_url);
let fields = vec!["project".to_string(), "updated".to_string()];
let tz = self.account_timezone().await?;
let max_pages = max_results.div_ceil(SEARCH_PAGE_SIZE) * 2 + 8;
let mut pager = KeysetPager::new(scope.since, max_pages);
let mut out: Vec<ChangelogIssue> = Vec::new();
let mut truncated = false;
loop {
let remaining = max_results.saturating_sub(out.len());
if remaining == 0 {
truncated = true;
break;
}
let page_size = remaining.min(SEARCH_PAGE_SIZE);
let request = pager.request();
let jql = build_jql(
&SyncScope {
project_key: scope.project_key.clone(),
since: request.since,
},
tz,
)?;
let body = json!({
"jql": jql,
"startAt": request.start_at,
"maxResults": page_size,
"fields": fields,
"expand": ["changelog"],
});
debug!(url = %url, %jql, start_at = request.start_at, "POST (with changelog)");
let parsed: ChangelogSearchResponse =
with_retry("search_with_changelog", &self.retry, &self.budget, || {
self.post(&url, &body)
})
.await?;
let issues: Vec<ChangelogIssue> = parsed
.issues
.into_iter()
.map(ChangelogIssue::from_api)
.collect();
let items: Vec<PagedItem> = issues.iter().map(|i| (i.key.clone(), i.updated)).collect();
let step = pager.record_page(&items, page_size);
for (issue, is_new) in issues.into_iter().zip(step.is_new) {
if !is_new {
continue;
}
out.push(issue);
if out.len() >= max_results {
break;
}
}
if !step.more {
break;
}
}
Ok(ChangelogWalk {
issues: out,
offset_paged_minute: pager.offset_paged_minute(),
truncated,
})
}
async fn post<T: serde::de::DeserializeOwned>(&self, url: &str, body: &Value) -> Result<T> {
post_json(&self.client, self.credentials.as_ref(), url, body).await
}
async fn get<T: serde::de::DeserializeOwned>(&self, url: &str) -> Result<T> {
get_json(&self.client, self.credentials.as_ref(), url).await
}
pub async fn fetch_comments(&self, key: &str) -> Result<Vec<JiraComment>> {
let mut out = Vec::new();
let mut start_at = 0u64;
for _ in 0..MAX_COMMENT_PAGES {
let url = format!(
"{}/rest/api/3/issue/{}/comment?startAt={}&maxResults={}",
self.base_url, key, start_at, COMMENT_PAGE_SIZE
);
debug!(url = %url, "GET");
let parsed: CommentSearchResponse =
with_retry("fetch_comments", &self.retry, &self.budget, || {
self.get(&url)
})
.await?;
let n = parsed.comments.len();
let page_size = parsed
.max_results
.map(|m| m.min(COMMENT_PAGE_SIZE))
.filter(|m| *m > 0);
for c in parsed.comments {
if let Some(comment) = JiraComment::from_api(c) {
out.push(comment);
}
}
let ended = match page_size {
Some(size) => n < size,
None => n == 0,
};
if ended {
return Ok(out);
}
start_at += n as u64;
}
Err(CollectError::PagingBudgetExceeded {
endpoint: "comment",
key: key.to_string(),
pages: MAX_COMMENT_PAGES,
})
}
}
#[cfg(test)]
#[path = "client_tests.rs"]
mod tests;