g1t/crates/runner/src/selfhosted/api.rs
| 1 | //! The runner's calls to g1t: outbound HTTPS only, JSON in `snake_case`. |
| 2 | //! See `g1t_contracts::runners` for the shapes. |
| 3 | |
| 4 | use std::time::Duration; |
| 5 | |
| 6 | use anyhow::{Result, anyhow}; |
| 7 | use serde::Deserialize; |
| 8 | use serde_json::{Map, Value, json}; |
| 9 | |
| 10 | use super::VERSION; |
| 11 | |
| 12 | #[derive(Clone, Debug, Deserialize)] |
| 13 | pub struct Assignment { |
| 14 | pub kind: String, |
| 15 | pub id: String, |
| 16 | pub name: String, |
| 17 | pub repo: String, |
| 18 | pub timeout_minutes: u32, |
| 19 | #[serde(default)] |
| 20 | pub image: Option<String>, |
| 21 | #[serde(default)] |
| 22 | pub token: Option<String>, |
| 23 | #[serde(default)] |
| 24 | pub env: Option<Map<String, Value>>, |
| 25 | } |
| 26 | |
| 27 | #[derive(Debug, Default, Deserialize)] |
| 28 | pub struct Poll { |
| 29 | #[serde(default)] |
| 30 | pub assignment: Option<Assignment>, |
| 31 | #[serde(default)] |
| 32 | pub cancel: Vec<String>, |
| 33 | #[serde(default)] |
| 34 | pub credential: Option<String>, |
| 35 | #[serde(default)] |
| 36 | pub removed: bool, |
| 37 | } |
| 38 | |
| 39 | #[derive(Debug, Deserialize)] |
| 40 | pub struct RegisteredRunner { |
| 41 | pub id: String, |
| 42 | pub name: String, |
| 43 | pub workspace: String, |
| 44 | #[serde(default)] |
| 45 | pub repo: Option<String>, |
| 46 | #[serde(default)] |
| 47 | pub group: Option<String>, |
| 48 | pub labels: Vec<String>, |
| 49 | } |
| 50 | |
| 51 | #[derive(Debug, Deserialize)] |
| 52 | pub struct Registered { |
| 53 | pub runner: RegisteredRunner, |
| 54 | pub credential: String, |
| 55 | } |
| 56 | |
| 57 | /// How a call failed: g1t refused it (with its status and message), or it |
| 58 | /// did not get through. |
| 59 | #[derive(Debug)] |
| 60 | pub enum Failure { |
| 61 | Refused(u16, String), |
| 62 | Unreachable(String), |
| 63 | } |
| 64 | |
| 65 | impl std::fmt::Display for Failure { |
| 66 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 67 | match self { |
| 68 | Failure::Refused(status, message) => write!(f, "{message} ({status})"), |
| 69 | Failure::Unreachable(why) => write!(f, "could not reach g1t: {why}"), |
| 70 | } |
| 71 | } |
| 72 | } |
| 73 | |
| 74 | fn agent() -> ureq::Agent { |
| 75 | ureq::AgentBuilder::new() |
| 76 | .timeout_connect(Duration::from_secs(15)) |
| 77 | .timeout_read(Duration::from_secs(60)) |
| 78 | .user_agent(&format!("g1t-runner/{VERSION} ({}; {})", std::env::consts::OS, std::env::consts::ARCH)) |
| 79 | .build() |
| 80 | } |
| 81 | |
| 82 | fn post(url: &str, credential: Option<&str>, body: &Value) -> std::result::Result<Value, Failure> { |
| 83 | let mut request = agent().post(url); |
| 84 | if let Some(credential) = credential { |
| 85 | request = request.set("authorization", &format!("Bearer {credential}")); |
| 86 | } |
| 87 | match request.send_json(body.clone()) { |
| 88 | Ok(response) => response.into_json::<Value>().map_err(|e| Failure::Unreachable(e.to_string())), |
| 89 | Err(ureq::Error::Status(status, response)) => { |
| 90 | let text = response.into_string().unwrap_or_default(); |
| 91 | let message = serde_json::from_str::<Value>(&text) |
| 92 | .ok() |
| 93 | .and_then(|v| v["error"]["message"].as_str().map(str::to_owned)) |
| 94 | .unwrap_or(text); |
| 95 | Err(Failure::Refused(status, message)) |
| 96 | } |
| 97 | Err(error) => Err(Failure::Unreachable(error.to_string())), |
| 98 | } |
| 99 | } |
| 100 | |
| 101 | pub struct Api { |
| 102 | pub base: String, |
| 103 | pub runner: String, |
| 104 | pub credential: String, |
| 105 | } |
| 106 | |
| 107 | impl Api { |
| 108 | pub fn register(base: &str, body: &Value) -> Result<Registered> { |
| 109 | let answer = post(&format!("{base}/runners/register"), None, body).map_err(|failure| anyhow!("{failure}"))?; |
| 110 | Ok(serde_json::from_value(answer)?) |
| 111 | } |
| 112 | |
| 113 | pub fn poll(&self, running: &[String], wait_ms: u64) -> std::result::Result<Poll, Failure> { |
| 114 | let body = json!({ "version": VERSION, "running": running, "wait_ms": wait_ms }); |
| 115 | let answer = post(&format!("{}/runners/{}/poll", self.base, self.runner), Some(&self.credential), &body)?; |
| 116 | serde_json::from_value(answer).map_err(|e| Failure::Unreachable(e.to_string())) |
| 117 | } |
| 118 | |
| 119 | pub fn finished(&self, id: &str, exit_code: i32, reason: Option<&str>) -> std::result::Result<(), Failure> { |
| 120 | let body = json!({ "id": id, "exit_code": exit_code, "reason": reason }); |
| 121 | // Tried a few times: a job that crashed should not wait for the |
| 122 | // silence check to be noticed. |
| 123 | let mut last = None; |
| 124 | for attempt in 0..3 { |
| 125 | match post(&format!("{}/runners/{}/finished", self.base, self.runner), Some(&self.credential), &body) { |
| 126 | Ok(_) => return Ok(()), |
| 127 | Err(Failure::Refused(status, message)) => return Err(Failure::Refused(status, message)), |
| 128 | Err(failure) => { |
| 129 | last = Some(failure); |
| 130 | std::thread::sleep(Duration::from_secs(2 * (attempt + 1))); |
| 131 | } |
| 132 | } |
| 133 | } |
| 134 | Err(last.unwrap_or(Failure::Unreachable("no answer".into()))) |
| 135 | } |
| 136 | |
| 137 | pub fn remove(&self) -> std::result::Result<(), Failure> { |
| 138 | post(&format!("{}/runners/{}/remove", self.base, self.runner), Some(&self.credential), &json!({})).map(|_| ()) |
| 139 | } |
| 140 | } |
| 141 | |
| 142 | /// Fetches a file's bytes. |
| 143 | pub fn download(url: &str) -> Result<Vec<u8>> { |
| 144 | let response = agent().get(url).timeout(Duration::from_secs(300)).call().map_err(|e| anyhow!("could not download {url}: {e}"))?; |
| 145 | let mut bytes = Vec::new(); |
| 146 | std::io::Read::read_to_end(&mut response.into_reader(), &mut bytes)?; |
| 147 | Ok(bytes) |
| 148 | } |