use chrono::{DateTime, NaiveDate, TimeZone, Utc};
use clap::Args;
use tracing::{info, warn};
use tga::collect::errors::CollectError;
use tga::collect::jira::sync::{
build_jql, plan_cursor, resolve_scope, validate_project_key, CursorPlan,
};
use tga::collect::jira::{ChangelogIssue, JiraClient, JiraTransition};
use tga::core::config::Config;
use tga::core::db::{
check_freshness, get_cursor, list_cursor_projects, set_cursor, upsert_comment_detail,
upsert_ticket_transition, CommentDetailRow, Database, FreshnessStatus, TicketTransitionRow,
};
const DEFAULT_MAX_TICKETS: usize = 10_000;
const MAX_CONSECUTIVE_TICKET_FAILURES: usize = 10;
#[derive(Args, Debug, Default)]
#[command(
about = "Sync JIRA status transitions and comments into fact_ticket_transitions / fact_jira_comment_detail.",
long_about = "Fetch JIRA changelog status transitions and full comment history for the\n\
configured (or --project-overridden) project, and persist them into\n\
`fact_ticket_transitions` and `fact_jira_comment_detail` in tga.db.\n\n\
Incremental by default: resumes from the stored `jira_sync_cursor` for the\n\
project. Pass --backfill for a full historical pull (first-ever sync of a\n\
project is always a full pull automatically, even without --backfill).\n\n\
Requires `jira.url` (and `jira.username`/`jira.token`, which may reference\n\
`${ENV_VAR}` placeholders) configured in config.yaml.\n\n\
JQL date literals carry no timezone and JIRA evaluates them in the querying\n\
account's profile timezone, so the sync window is rendered in that zone. It\n\
is read from `GET /rest/api/3/myself`; set `jira.timezone` (an IANA name such\n\
as `UTC`) to pin it explicitly when that endpoint is not reachable. The sync\n\
refuses to run rather than guessing.",
after_help = "EXAMPLES:\n\
# Incremental sync using the stored cursor (or full history on first run)\n\
tga jira sync\n\n\
# Full historical backfill, ignoring any stored cursor\n\
tga jira sync --backfill\n\n\
# Sync only tickets updated on/after a specific date\n\
tga jira sync --since 2026-01-01\n\n\
# Preview without writing to the database\n\
tga jira sync --dry-run\n\n\
TIPS:\n\
- Run `tga jira freshness` after a sync (or on a schedule) to catch a\n\
silently-stopped sync before downstream reports serve stale data."
)]
pub struct JiraSyncArgs {
#[arg(long, value_name = "KEY")]
pub project: Option<String>,
#[arg(long, value_name = "DATE")]
pub since: Option<String>,
#[arg(long, default_value_t = false)]
pub backfill: bool,
#[arg(long, value_name = "N")]
pub max_tickets: Option<usize>,
#[arg(long, default_value_t = false)]
pub dry_run: bool,
}
#[derive(Args, Debug)]
#[command(
about = "Check freshness of the JIRA-derived fact tables (fails loudly if stale/empty).",
long_about = "Report row counts and sync recency for `fact_ticket_transitions` and\n\
`fact_jira_comment_detail`. Exits non-zero (unless --report-only) if either\n\
table is empty or has not been written to within --max-age-days.\n\n\
This is the CRITICAL freshness guard from issue #3966: intended to be run\n\
as a health check (e.g. from the same cron slot as `tga jira sync`, or a\n\
separate monitoring job) so a sync that silently stopped running is caught\n\
loudly instead of downstream reports serving stale data with no alarm.",
after_help = "EXAMPLES:\n\
# Standard health check: every project with a sync cursor, checked\n\
# individually (fails the process if ANY project's table is stale)\n\
tga jira freshness\n\n\
# Check one project only\n\
tga jira freshness --project PROJ\n\n\
# Report only, never fail the process (e.g. informational dashboard use)\n\
tga jira freshness --report-only --max-age-days 7"
)]
pub struct JiraFreshnessArgs {
#[arg(long, default_value_t = 2)]
pub max_age_days: i64,
#[arg(long, default_value_t = false)]
pub report_only: bool,
#[arg(long, value_name = "KEY")]
pub project: Option<String>,
#[arg(long, value_name = "DAYS")]
pub max_cursor_lag_days: Option<i64>,
}
fn parse_cli_date(s: &str) -> anyhow::Result<DateTime<Utc>> {
let d = NaiveDate::parse_from_str(s, "%Y-%m-%d")
.map_err(|e| anyhow::anyhow!("invalid --since date '{s}' (expected YYYY-MM-DD): {e}"))?;
let ndt = d
.and_hms_opt(0, 0, 0)
.ok_or_else(|| anyhow::anyhow!("invalid time-of-day for date '{s}'"))?;
Ok(Utc.from_utc_datetime(&ndt))
}
fn resolve_project_key(config: &Config, cli_project: Option<&str>) -> anyhow::Result<String> {
let key = match cli_project {
Some(p) => p.to_string(),
None => config
.jira
.as_ref()
.and_then(|j| j.project_key.clone())
.filter(|p| !p.is_empty())
.ok_or_else(|| {
anyhow::anyhow!(
"no JIRA project scope: pass --project <KEY> or set jira.project_key in config.yaml"
)
})?,
};
validate_project_key(&key).map_err(|e| anyhow::anyhow!(e))?;
Ok(key)
}
fn build_client(config: &Config) -> anyhow::Result<JiraClient> {
let jira_config = config
.jira
.clone()
.ok_or_else(|| anyhow::anyhow!("`jira:` section is missing from config.yaml"))?;
JiraClient::new(&jira_config).map_err(|e| match e {
CollectError::Config(msg) => anyhow::anyhow!("{msg}"),
other => anyhow::anyhow!(other),
})
}
pub async fn run_sync(config: Config, db: &mut Database, args: JiraSyncArgs) -> anyhow::Result<()> {
let project_key = resolve_project_key(&config, args.project.as_deref())?;
let client = build_client(&config)?;
let explicit_since = args.since.as_deref().map(parse_cli_date).transpose()?;
let stored_cursor = get_cursor(db.connection(), &project_key)?
.and_then(|c| DateTime::parse_from_rfc3339(&c.last_synced_at).ok())
.map(|d| d.with_timezone(&Utc));
let scope = resolve_scope(&project_key, explicit_since, args.backfill, stored_cursor);
let max_tickets = args.max_tickets.unwrap_or(DEFAULT_MAX_TICKETS);
let tz = client.account_timezone().await?;
let logged_jql = build_jql(&scope, tz)?;
info!(
project = %project_key,
jql = %logged_jql,
timezone = %tz,
backfill = args.backfill,
dry_run = args.dry_run,
"starting tga jira sync"
);
let walk = client.search_with_changelog(&scope, max_tickets).await?;
let mut tickets_scanned = 0usize;
let mut transitions_written = 0usize;
let mut comments_ingested = 0usize;
let mut observed_updated: Vec<DateTime<Utc>> = Vec::new();
let mut failed_tickets: Vec<String> = Vec::new();
let mut failed_updated: Vec<Option<DateTime<Utc>>> = Vec::new();
let mut consecutive_failures = 0usize;
let mut tripped = false;
for issue in &walk.issues {
tickets_scanned += 1;
if let Some(u) = issue.updated {
observed_updated.push(u);
}
let repaired = match issue.truncated_history_total {
Some(expected) => match client.fetch_changelog(&issue.key, Some(expected)).await {
Ok(full) => Some(full),
Err(e) => {
warn!(
ticket = %issue.key,
error = %e,
"could not repair this ticket's truncated changelog; skipping its \
transitions rather than persisting a knowingly-short history, and \
holding the sync cursor at or below it"
);
record_failure(
issue,
&mut failed_tickets,
&mut failed_updated,
&mut consecutive_failures,
);
if consecutive_failures >= MAX_CONSECUTIVE_TICKET_FAILURES {
tripped = true;
break;
}
continue;
}
},
None => None,
};
let transitions = repaired.as_deref().unwrap_or(&issue.transitions);
if !args.dry_run {
write_transitions(db, issue, transitions)?;
}
transitions_written += transitions.len();
match client.fetch_comments(&issue.key).await {
Ok(comments) => {
consecutive_failures = 0;
comments_ingested += comments.len();
if args.dry_run {
continue;
}
for c in &comments {
let row = CommentDetailRow {
ticket_key: issue.key.clone(),
comment_id: c.id.clone(),
project_key: issue.project_key.clone(),
author: c.author.clone(),
created_at: c.created.to_rfc3339(),
body_len: c.body_len,
};
upsert_comment_detail(db.connection(), &row)?;
}
}
Err(e) => {
warn!(
ticket = %issue.key,
error = %e,
"failed to fetch comments for this ticket; holding the sync \
cursor at or below it so the next run re-fetches it"
);
record_failure(
issue,
&mut failed_tickets,
&mut failed_updated,
&mut consecutive_failures,
);
if consecutive_failures >= MAX_CONSECUTIVE_TICKET_FAILURES {
tripped = true;
break;
}
}
}
}
if tripped {
warn!(
consecutive_failures,
tickets_scanned,
"aborting the walk: too many consecutive per-ticket failures. Continuing \
would issue thousands more requests against a remote that is already \
failing every one of them"
);
}
if !args.dry_run {
match plan_cursor(&observed_updated, &failed_updated, stored_cursor) {
CursorPlan::Advance(next) => set_cursor(
db.connection(),
&project_key,
&next.to_rfc3339(),
tickets_scanned as i64,
)?,
CursorPlan::Hold => info!(
project = %project_key,
tickets_scanned,
failures = failed_tickets.len(),
"cursor left unchanged (no usable `updated` timestamps, or a \
failed ticket could not be placed on the timeline)"
),
}
}
println!(
"JIRA sync ({project_key}): {tickets_scanned} ticket(s) scanned, \
{transitions_written} transition(s), {comments_ingested} comment(s), \
{} failed ticket(s){}.",
failed_tickets.len(),
if args.dry_run {
" [dry-run: no writes]"
} else {
""
}
);
if walk.truncated {
println!(
" note: stopped at the --max-tickets limit ({max_tickets}); more tickets match \
this window. Re-run to continue from the recorded cursor."
);
warn!(
project = %project_key,
max_tickets,
"changelog walk truncated at --max-tickets; the window is only partly covered"
);
}
if let Some(minute) = walk.offset_paged_minute {
println!(
" note: {} contained more tickets than one page, so that minute was walked by \
offset. A ticket edited during that walk could have been missed; re-cover it \
with `tga jira sync --project {project_key} --since {}` if in doubt.",
minute.to_rfc3339(),
minute.date_naive()
);
warn!(
project = %project_key,
minute = %minute.to_rfc3339(),
"walked a single minute by offset; see `collect::jira::paging` for the residual"
);
}
if !failed_tickets.is_empty() {
anyhow::bail!(
"JIRA sync ({project_key}) could not fully ingest {} of {} ticket(s): {}.{} \
The sync cursor was held at or below the earliest failure, so the next run \
re-fetches them — but this run's data is incomplete and must not be treated \
as a successful sync.",
failed_tickets.len(),
tickets_scanned,
summarize_keys(&failed_tickets),
if tripped {
format!(
" The walk was ABORTED after {MAX_CONSECUTIVE_TICKET_FAILURES} consecutive \
failures rather than continuing through the remaining tickets — the remote \
is likely rate-limiting or down, so retry later rather than immediately."
)
} else {
String::new()
},
);
}
Ok(())
}
fn summarize_keys(keys: &[String]) -> String {
const MAX_LISTED: usize = 10;
if keys.len() <= MAX_LISTED {
return keys.join(", ");
}
format!(
"{}, … and {} more",
keys[..MAX_LISTED].join(", "),
keys.len() - MAX_LISTED
)
}
fn record_failure(
issue: &ChangelogIssue,
failed_tickets: &mut Vec<String>,
failed_updated: &mut Vec<Option<DateTime<Utc>>>,
consecutive_failures: &mut usize,
) {
failed_tickets.push(issue.key.clone());
failed_updated.push(issue.updated);
*consecutive_failures += 1;
}
fn write_transitions(
db: &mut Database,
issue: &ChangelogIssue,
transitions: &[JiraTransition],
) -> anyhow::Result<()> {
let conn = db.connection_mut();
let tx = conn.transaction()?;
for t in transitions {
let row = TicketTransitionRow {
ticket_key: issue.key.clone(),
project_key: issue.project_key.clone(),
from_status: t.from_status.clone(),
to_status: t.to_status.clone(),
transitioned_at: t.created.to_rfc3339(),
author: t.author.clone(),
};
upsert_ticket_transition(&tx, &row)?;
}
tx.commit()?;
Ok(())
}
pub fn run_freshness(
config: &Config,
db: &Database,
args: JiraFreshnessArgs,
) -> anyhow::Result<()> {
let scopes: Vec<Option<String>> = match &args.project {
Some(p) => {
validate_project_key(p).map_err(|e| anyhow::anyhow!(e))?;
vec![Some(p.clone())]
}
None => {
let mut projects = list_cursor_projects(db.connection())?;
if let Some(configured) = config
.jira
.as_ref()
.and_then(|j| j.project_key.clone())
.filter(|p| !p.is_empty())
{
if !projects.contains(&configured) {
projects.push(configured);
}
}
projects.sort();
if projects.is_empty() {
vec![None]
} else {
projects.into_iter().map(Some).collect()
}
}
};
let mut statuses: Vec<FreshnessStatus> = Vec::new();
for scope in &scopes {
statuses.extend(check_freshness(
db.connection(),
args.max_age_days,
scope.as_deref(),
)?);
}
let mut any_stale = false;
for s in &statuses {
let age_desc = match s.age_seconds {
Some(age) => format!("{:.1}d old", age as f64 / 86_400.0),
None => "no rows".to_string(),
};
let verdict = if s.stale { "STALE" } else { "OK" };
println!(
"{:<10} {:<28} rows={:<8} last_synced={:<14} [{}]",
s.project.as_deref().unwrap_or("<all>"),
s.table,
s.row_count,
age_desc,
verdict
);
if s.stale {
any_stale = true;
}
}
let mut lagging: Vec<String> = Vec::new();
for project in scopes.iter().flatten() {
let Some(cursor) = get_cursor(db.connection(), project)? else {
continue;
};
let Ok(parsed) = DateTime::parse_from_rfc3339(&cursor.last_synced_at) else {
continue;
};
let lag_days = (Utc::now() - parsed.with_timezone(&Utc)).num_seconds() as f64 / 86_400.0;
let flagged = args
.max_cursor_lag_days
.is_some_and(|max| lag_days > max as f64);
println!(
"{project:<10} {:<28} cursor={} ({lag_days:.1}d behind){}",
"jira_sync_cursor",
cursor.last_synced_at,
if flagged { " [LAGGING]" } else { "" }
);
if flagged {
lagging.push(project.clone());
}
}
if !lagging.is_empty() {
any_stale = true;
}
if any_stale {
let msg = format!(
"one or more JIRA fact tables are stale or empty (threshold: {} day(s), \
scopes checked: {}); see rows above",
args.max_age_days,
scopes
.iter()
.map(|s| s.as_deref().unwrap_or("<all>"))
.collect::<Vec<_>>()
.join(", ")
);
if args.report_only {
warn!("{msg}");
return Ok(());
}
anyhow::bail!(msg);
}
Ok(())
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;