use crate::config::Config;
use crate::model::{
ArtifactNode, DashboardEmbedded, DashboardResponse, HistoryResponse, PipelineInstance,
ViewFilters,
};
use anyhow::{Context, Result};
use reqwest::Method;
use reqwest::blocking::{Client, RequestBuilder};
#[derive(Clone)]
pub struct GoCdClient {
client: Client,
base_url: String,
username: Option<String>,
password: Option<String>,
token: Option<String>,
}
impl GoCdClient {
pub fn new(cfg: &Config) -> Result<Self> {
let client = Client::builder()
.danger_accept_invalid_certs(cfg.insecure_skip_verify)
.timeout(std::time::Duration::from_secs(45))
.connect_timeout(std::time::Duration::from_secs(8))
.build()
.context("building HTTP client")?;
Ok(GoCdClient {
client,
base_url: cfg.server_url.trim_end_matches('/').to_string(),
username: cfg.username.clone(),
password: cfg.password.clone(),
token: cfg.auth_token.clone(),
})
}
fn request(&self, method: Method, path: &str, api_version: u8) -> RequestBuilder {
let url = format!("{}{}", self.base_url, path);
let mut rb = self.client.request(method, url).header(
"Accept",
format!("application/vnd.go.cd.v{api_version}+json"),
);
if let Some(token) = &self.token {
rb = rb.bearer_auth(token);
} else if let Some(user) = &self.username {
rb = rb.basic_auth(user, self.password.clone());
}
rb
}
fn request_raw(&self, method: Method, path: &str) -> RequestBuilder {
let url = format!("{}{}", self.base_url, path);
let mut rb = self.client.request(method, url);
if let Some(token) = &self.token {
rb = rb.bearer_auth(token);
} else if let Some(user) = &self.username {
rb = rb.basic_auth(user, self.password.clone());
}
rb
}
pub fn fetch_dashboard(
&self,
etag: Option<&str>,
view: Option<&str>,
) -> Result<Option<(DashboardEmbedded, Option<String>)>> {
let mut rb = self.request(Method::GET, "/api/dashboard", 4);
if let Some(name) = view {
rb = rb.query(&[("viewName", name)]);
}
if let Some(tag) = etag {
rb = rb.header("If-None-Match", tag);
}
let resp = rb.send().context("requesting dashboard")?;
let status = resp.status();
if status == reqwest::StatusCode::NOT_MODIFIED {
return Ok(None);
}
let new_etag = resp
.headers()
.get(reqwest::header::ETAG)
.and_then(|v| v.to_str().ok())
.map(str::to_string);
let body = resp.text().context("reading dashboard response body")?;
if !status.is_success() {
anyhow::bail!("GoCD returned {status} for dashboard: {}", truncate(&body));
}
let parsed: DashboardResponse = serde_json::from_str(&body)
.with_context(|| format!("parsing dashboard response: {}", truncate(&body)))?;
Ok(Some((parsed.embedded, new_etag)))
}
pub fn fetch_history_page(
&self,
pipeline_name: &str,
after: Option<u64>,
) -> Result<(Vec<PipelineInstance>, Option<u64>)> {
let mut path = format!("/api/pipelines/{pipeline_name}/history");
if let Some(cursor) = after {
path.push_str(&format!("?after={cursor}"));
}
let resp = self
.request(Method::GET, &path, 1)
.send()
.with_context(|| format!("requesting history for {pipeline_name}"))?;
let status = resp.status();
let body = resp.text().context("reading history response body")?;
if !status.is_success() {
anyhow::bail!(
"GoCD returned {status} for {pipeline_name} history: {}",
truncate(&body)
);
}
let parsed: HistoryResponse = serde_json::from_str(&body).with_context(|| {
format!(
"parsing history response for {pipeline_name}: {}",
truncate(&body)
)
})?;
let next = parsed
.links
.and_then(|l| l.next)
.and_then(|n| crate::model::next_page_cursor(&n.href));
Ok((parsed.pipelines, next))
}
pub fn trigger_pipeline(&self, pipeline_name: &str) -> Result<()> {
let path = format!("/api/pipelines/{pipeline_name}/schedule");
let resp = self
.request(Method::POST, &path, 1)
.header("X-GoCD-Confirm", "true")
.header("Content-Type", "application/json")
.body("{}")
.send()
.with_context(|| format!("triggering {pipeline_name}"))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} triggering {pipeline_name}: {}",
truncate(&body)
);
}
Ok(())
}
pub fn trigger_pipeline_with_vars(
&self,
pipeline_name: &str,
vars: &[(String, String)],
) -> Result<()> {
let env: Vec<serde_json::Value> = vars
.iter()
.map(|(name, value)| serde_json::json!({ "name": name, "value": value, "secure": false }))
.collect();
let path = format!("/api/pipelines/{pipeline_name}/schedule");
let resp = self
.request(Method::POST, &path, 1)
.header("X-GoCD-Confirm", "true")
.json(&serde_json::json!({
"environment_variables": env,
"update_materials_before_scheduling": true,
}))
.send()
.with_context(|| format!("triggering {pipeline_name} with variables"))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} triggering {pipeline_name}: {}",
truncate(&body)
);
}
Ok(())
}
pub fn rerun_failed_jobs(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
) -> Result<()> {
self.rerun(
pipeline_name,
pipeline_counter,
stage_name,
stage_counter,
"run-failed-jobs",
)
}
pub fn rerun_stage(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
) -> Result<()> {
self.rerun(
pipeline_name,
pipeline_counter,
stage_name,
stage_counter,
"run",
)
}
fn rerun(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
verb: &str,
) -> Result<()> {
let path = format!(
"/api/stages/{pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}/{verb}"
);
let resp = self
.request(Method::POST, &path, 3)
.header("X-GoCD-Confirm", "true")
.send()
.with_context(|| format!("rerunning ({verb}) {pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}"))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} rerunning stage: {}",
truncate(&body)
);
}
Ok(())
}
pub fn pause_pipeline(&self, pipeline_name: &str, cause: &str) -> Result<()> {
let path = format!("/api/pipelines/{pipeline_name}/pause");
let resp = self
.request(Method::POST, &path, 1)
.header("X-GoCD-Confirm", "true")
.json(&serde_json::json!({ "pause_cause": cause }))
.send()
.with_context(|| format!("pausing {pipeline_name}"))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} pausing {pipeline_name}: {}",
truncate(&body)
);
}
Ok(())
}
pub fn unpause_pipeline(&self, pipeline_name: &str) -> Result<()> {
let path = format!("/api/pipelines/{pipeline_name}/unpause");
let resp = self
.request(Method::POST, &path, 1)
.header("X-GoCD-Confirm", "true")
.send()
.with_context(|| format!("unpausing {pipeline_name}"))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} unpausing {pipeline_name}: {}",
truncate(&body)
);
}
Ok(())
}
pub fn cancel_stage(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
) -> Result<()> {
let path = format!(
"/api/stages/{pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}/cancel"
);
let resp = self
.request(Method::POST, &path, 3)
.header("X-GoCD-Confirm", "true")
.send()
.with_context(|| {
format!(
"cancelling {pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}"
)
})?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!(
"GoCD returned {status} cancelling stage: {}",
truncate(&body)
);
}
Ok(())
}
pub fn fetch_views(&self) -> Result<(ViewFilters, Option<String>)> {
let resp = self
.request(Method::GET, "/api/internal/pipeline_selection", 1)
.send()
.context("requesting personalized views")?;
let status = resp.status();
let etag = resp
.headers()
.get(reqwest::header::ETAG)
.and_then(|v| v.to_str().ok())
.map(|t| t.replace("--gzip", ""));
let body = resp.text().context("reading views body")?;
if !status.is_success() {
anyhow::bail!(
"GoCD returned {status} for pipeline_selection: {}",
truncate(&body)
);
}
let filters = serde_json::from_str(&body)
.with_context(|| format!("parsing views: {}", truncate(&body)))?;
Ok((filters, etag))
}
pub fn save_view(&self, name: &str, pipelines: Vec<String>) -> Result<()> {
let (mut current, etag) = self.fetch_views()?;
let new_filter = crate::model::ViewFilter {
name: name.to_string(),
kind: "whitelist".to_string(),
state: Vec::new(),
pipelines,
};
match current.filters.iter_mut().find(|f| f.name == name) {
Some(existing) => *existing = new_filter,
None => current.filters.push(new_filter),
}
let mut rb = self
.request(Method::PUT, "/api/internal/pipeline_selection", 1)
.header("Content-Type", "application/json")
.json(¤t);
if let Some(tag) = etag {
rb = rb.header("If-Match", tag);
}
let resp = rb.send().context("saving view")?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().unwrap_or_default();
anyhow::bail!("GoCD returned {status} saving view: {}", truncate(&body));
}
Ok(())
}
pub fn fetch_artifacts(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
job_name: &str,
) -> Result<Vec<ArtifactNode>> {
let path = format!(
"/files/{pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}/{job_name}.json"
);
let resp = self
.request_raw(Method::GET, &path)
.send()
.context("requesting artifacts")?;
let status = resp.status();
let body = resp.text().context("reading artifacts body")?;
if !status.is_success() {
anyhow::bail!("GoCD returned {status} for artifacts: {}", truncate(&body));
}
serde_json::from_str(&body)
.with_context(|| format!("parsing artifacts: {}", truncate(&body)))
}
pub fn fetch_console_log(
&self,
pipeline_name: &str,
pipeline_counter: i64,
stage_name: &str,
stage_counter: &str,
job_name: &str,
start_line: usize,
) -> Result<String> {
let mut path = format!(
"/files/{pipeline_name}/{pipeline_counter}/{stage_name}/{stage_counter}/{job_name}/cruise-output/console.log"
);
if start_line > 0 {
path.push_str(&format!("?startLineNumber={start_line}"));
}
let resp = self
.request_raw(Method::GET, &path)
.send()
.context("requesting console log")?;
let status = resp.status();
let body = resp.text().context("reading console log body")?;
if !status.is_success() {
anyhow::bail!(
"GoCD returned {status} for console log: {}",
truncate(&body)
);
}
Ok(body)
}
}
fn truncate(s: &str) -> String {
let t = s.trim();
if t.starts_with('<') || t.to_ascii_lowercase().contains("<html") {
return "a proxy or gateway answered instead of GoCD".to_string();
}
t.chars().take(300).collect()
}