use std::path::Path;
use anyhow::{anyhow, Context};
use chrono::{DateTime, Utc};
use rusqlite::params;
use serde::Deserialize;
use tga::core::db::Database;
use super::CollectStats;
#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
pub(super) struct FileDeploymentRecord {
pub deploy_id: String,
pub repo: String,
#[serde(default)]
pub environment: Option<String>,
pub triggered_at: String,
#[serde(default)]
pub completed_at: Option<String>,
pub status: String,
#[serde(default)]
pub git_sha: Option<String>,
#[serde(default)]
pub git_tag: Option<String>,
#[serde(default)]
pub triggered_by_pr: Option<i64>,
#[serde(default)]
pub source: Option<String>,
}
#[derive(Debug, Deserialize)]
struct FileDeploymentsDoc {
deployments: Vec<FileDeploymentRecord>,
}
pub(super) fn ingest_file(db: &mut Database, path: &Path) -> anyhow::Result<CollectStats> {
let records = load_records(path)?;
let mut stats = CollectStats::default();
let conn = db.connection_mut();
let tx = conn.transaction()?;
{
let mut insert = tx.prepare(
"INSERT OR IGNORE INTO fact_deployments \
(deploy_id, repo, environment, triggered_at, completed_at, \
status, git_sha, git_tag, triggered_by_pr, source) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
)?;
for (idx, rec) in records.iter().enumerate() {
let row_ctx = format!("record #{} (deploy_id={:?})", idx + 1, rec.deploy_id);
if rec.deploy_id.trim().is_empty() {
anyhow::bail!("{row_ctx}: deploy_id must not be empty");
}
if rec.repo.trim().is_empty() {
anyhow::bail!("{row_ctx}: repo must not be empty");
}
if rec.status.trim().is_empty() {
anyhow::bail!("{row_ctx}: status must not be empty");
}
let triggered_at = parse_rfc3339(&rec.triggered_at, "triggered_at", &row_ctx)?;
let completed_at = rec
.completed_at
.as_deref()
.filter(|s| !s.trim().is_empty())
.map(|s| parse_rfc3339(s, "completed_at", &row_ctx))
.transpose()?;
let environment = rec
.environment
.as_deref()
.filter(|s| !s.trim().is_empty())
.unwrap_or("production");
let source = rec
.source
.as_deref()
.filter(|s| !s.trim().is_empty())
.unwrap_or("file");
let changed = insert.execute(params![
rec.deploy_id,
rec.repo,
environment,
triggered_at.to_rfc3339(),
completed_at.map(|d: DateTime<Utc>| d.to_rfc3339()),
rec.status,
rec.git_sha,
rec.git_tag,
rec.triggered_by_pr,
source,
])?;
if changed > 0 {
stats.inserted += 1;
} else {
stats.skipped += 1;
}
}
}
tx.commit()?;
Ok(stats)
}
fn parse_rfc3339(s: &str, field: &str, row_ctx: &str) -> anyhow::Result<DateTime<Utc>> {
DateTime::parse_from_rfc3339(s)
.map(|d| d.with_timezone(&Utc))
.map_err(|e| anyhow!("{row_ctx}: {field} {s:?} is not a valid RFC3339 timestamp: {e}"))
}
fn load_records(path: &Path) -> anyhow::Result<Vec<FileDeploymentRecord>> {
match path.extension().and_then(|e| e.to_str()) {
Some("json") => load_json(path),
Some("csv") => load_csv(path),
other => anyhow::bail!(
"unsupported deployments file extension {other:?} for {}: expected .json or .csv",
path.display()
),
}
}
fn load_json(path: &Path) -> anyhow::Result<Vec<FileDeploymentRecord>> {
let body = std::fs::read_to_string(path)
.with_context(|| format!("failed to read deployments file: {}", path.display()))?;
let doc: FileDeploymentsDoc = serde_json::from_str(&body).with_context(|| {
format!(
"failed to parse {} as a deployments JSON document \
(expected {{\"deployments\": [...]}})",
path.display()
)
})?;
Ok(doc.deployments)
}
const CSV_COLUMNS: [&str; 10] = [
"deploy_id",
"repo",
"environment",
"triggered_at",
"completed_at",
"status",
"git_sha",
"git_tag",
"triggered_by_pr",
"source",
];
fn load_csv(path: &Path) -> anyhow::Result<Vec<FileDeploymentRecord>> {
let mut rdr = csv::ReaderBuilder::new()
.has_headers(true)
.from_path(path)
.with_context(|| format!("failed to open deployments CSV file: {}", path.display()))?;
let headers = rdr.headers()?.clone();
let mut col_idx = std::collections::HashMap::new();
for name in CSV_COLUMNS {
if let Some(pos) = headers.iter().position(|h| h == name) {
col_idx.insert(name, pos);
}
}
for required in ["deploy_id", "repo", "triggered_at", "status"] {
if !col_idx.contains_key(required) {
anyhow::bail!(
"{}: CSV header is missing required column '{required}' \
(expected columns: {})",
path.display(),
CSV_COLUMNS.join(", ")
);
}
}
let get = |record: &csv::StringRecord, col: &str| -> Option<String> {
col_idx
.get(col)
.and_then(|&i| record.get(i))
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
};
let mut out = Vec::new();
for (row_num, result) in rdr.records().enumerate() {
let record = result.with_context(|| {
format!(
"{}: malformed CSV at data row {}",
path.display(),
row_num + 1
)
})?;
let triggered_by_pr = match get(&record, "triggered_by_pr") {
Some(s) => Some(s.parse::<i64>().map_err(|e| {
anyhow!(
"{}: data row {}: triggered_by_pr {s:?} is not an integer: {e}",
path.display(),
row_num + 1
)
})?),
None => None,
};
out.push(FileDeploymentRecord {
deploy_id: get(&record, "deploy_id").unwrap_or_default(),
repo: get(&record, "repo").unwrap_or_default(),
environment: get(&record, "environment"),
triggered_at: get(&record, "triggered_at").unwrap_or_default(),
completed_at: get(&record, "completed_at"),
status: get(&record, "status").unwrap_or_default(),
git_sha: get(&record, "git_sha"),
git_tag: get(&record, "git_tag"),
triggered_by_pr,
source: get(&record, "source"),
});
}
Ok(out)
}
#[cfg(test)]
#[path = "file_source_tests.rs"]
mod tests;