g1t/crates/runner/src/selfhosted/api.rs
Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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 | } |