| 1 | //! The dial-out half: claim a job, fetch its checkout, report what happened. |
| 2 | //! |
| 3 | //! Every request originates here. anvil never connects to a runner, which is |
| 4 | //! the whole reason this is pull-based — the runner sits on a home network |
| 5 | //! behind NAT, and the forge sits on a public VPS. An inbound path from the |
| 6 | //! latter to the former would make the forge's compromise the build host's. |
| 7 | |
| 8 | use std::time::Duration; |
| 9 | |
| 10 | use anvil_job::{ |
| 11 | ArtifactSpec, |
| 12 | JobResult, |
| 13 | JobSpec, |
| 14 | RunnerInfo, |
| 15 | Stored, |
| 16 | }; |
| 17 | |
| 18 | use crate::executor::ArtifactSink; |
| 19 | |
| 20 | /// Generous relative to the server's claim poll (55s): a claim that finds work |
| 21 | /// returns at once, and one that does not still has to come back under this. |
| 22 | const REQUEST_TIMEOUT: Duration = Duration::from_secs(90); |
| 23 | |
| 24 | pub struct Client { |
| 25 | http: reqwest::Client, |
| 26 | base: String, |
| 27 | token: String, |
| 28 | info: RunnerInfo, |
| 29 | } |
| 30 | |
| 31 | impl Client { |
| 32 | /// Build the client every request goes through. |
| 33 | /// |
| 34 | /// The rustls *ring* provider is installed process-wide first: the |
| 35 | /// workspace builds reqwest with `-no-provider` on purpose (aws-lc-rs |
| 36 | /// needs cmake and breaks the zig musl cross-build), so a process that |
| 37 | /// does not install one panics inside `build()` — before it has claimed |
| 38 | /// anything, and identically on every host. |
| 39 | pub fn new(base: &str, token: &str, info: RunnerInfo) -> Result<Self, String> { |
| 40 | static TLS_PROVIDER: std::sync::Once = std::sync::Once::new(); |
| 41 | TLS_PROVIDER.call_once(|| { |
| 42 | let _ = rustls::crypto::ring::default_provider().install_default(); |
| 43 | }); |
| 44 | |
| 45 | let http = reqwest::Client::builder() |
| 46 | .timeout(REQUEST_TIMEOUT) |
| 47 | .build() |
| 48 | .map_err(|e| format!("building http client: {e}"))?; |
| 49 | Ok(Self { |
| 50 | http, |
| 51 | base: base.trim_end_matches('/').to_string(), |
| 52 | token: token.to_string(), |
| 53 | info, |
| 54 | }) |
| 55 | } |
| 56 | |
| 57 | fn url(&self, path: &str) -> String { |
| 58 | format!("{}{path}", self.base) |
| 59 | } |
| 60 | |
| 61 | /// Token plus identity on every request. The name goes in a header rather |
| 62 | /// than only in the claim body so the endpoints with no body of their own |
| 63 | /// (the checkout GET, artifact uploads) can be authorized the same way. |
| 64 | fn auth(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder { |
| 65 | req.header("X-Anvil-Runner-Token", &self.token) |
| 66 | .header("X-Anvil-Runner-Name", &self.info.name) |
| 67 | } |
| 68 | |
| 69 | /// Long-poll for a job. `Ok(None)` means the poll expired with nothing |
| 70 | /// queued, which is the common case and not an error. |
| 71 | pub async fn claim(&self) -> Result<Option<JobSpec>, String> { |
| 72 | let resp = self |
| 73 | .auth(self.http.post(self.url("/-/runner/claim"))) |
| 74 | .json(&self.info) |
| 75 | .send() |
| 76 | .await |
| 77 | .map_err(|e| format!("claim: {e}"))?; |
| 78 | match resp.status() { |
| 79 | reqwest::StatusCode::NO_CONTENT => Ok(None), |
| 80 | s if s.is_success() => resp |
| 81 | .json::<JobSpec>() |
| 82 | .await |
| 83 | .map(Some) |
| 84 | .map_err(|e| format!("decoding job: {e}")), |
| 85 | s => Err(format!("claim: {s} {}", body_hint(resp).await)), |
| 86 | } |
| 87 | } |
| 88 | |
| 89 | pub async fn checkout(&self, run_id: i64) -> Result<Vec<u8>, String> { |
| 90 | let resp = self |
| 91 | .auth( |
| 92 | self.http |
| 93 | .get(self.url(&format!("/-/runner/jobs/{run_id}/checkout.tar"))), |
| 94 | ) |
| 95 | .send() |
| 96 | .await |
| 97 | .map_err(|e| format!("checkout: {e}"))?; |
| 98 | if !resp.status().is_success() { |
| 99 | return Err(format!( |
| 100 | "checkout: {} {}", |
| 101 | resp.status(), |
| 102 | body_hint(resp).await |
| 103 | )); |
| 104 | } |
| 105 | resp.bytes() |
| 106 | .await |
| 107 | .map(|b| b.to_vec()) |
| 108 | .map_err(|e| format!("reading checkout: {e}")) |
| 109 | } |
| 110 | |
| 111 | /// Extend the lease. `false` means anvil no longer thinks we hold this run |
| 112 | /// — it was requeued out from under us, and the job should be abandoned |
| 113 | /// rather than reported. |
| 114 | pub async fn heartbeat(&self, run_id: i64) -> bool { |
| 115 | let sent = self |
| 116 | .auth( |
| 117 | self.http |
| 118 | .post(self.url(&format!("/-/runner/jobs/{run_id}/heartbeat"))) |
| 119 | .json(&self.info), |
| 120 | ) |
| 121 | .send() |
| 122 | .await; |
| 123 | matches!(sent, Ok(r) if r.status().is_success()) |
| 124 | } |
| 125 | |
| 126 | pub async fn report(&self, run_id: i64, result: &JobResult) -> Result<(), String> { |
| 127 | let resp = self |
| 128 | .auth( |
| 129 | self.http |
| 130 | .post(self.url(&format!("/-/runner/jobs/{run_id}/result"))) |
| 131 | .json(result), |
| 132 | ) |
| 133 | .send() |
| 134 | .await |
| 135 | .map_err(|e| format!("result: {e}"))?; |
| 136 | if resp.status().is_success() { |
| 137 | Ok(()) |
| 138 | } else { |
| 139 | Err(format!( |
| 140 | "result: {} {}", |
| 141 | resp.status(), |
| 142 | body_hint(resp).await |
| 143 | )) |
| 144 | } |
| 145 | } |
| 146 | } |
| 147 | |
| 148 | /// Uploads each artifact tar to anvil as it comes out of the container. |
| 149 | /// |
| 150 | /// The server does the storing, and reports back what it wrote — the runner |
| 151 | /// never learns anvil's on-disk layout, and `browse` stays the server's call. |
| 152 | pub struct UploadSink<'a> { |
| 153 | pub client: &'a Client, |
| 154 | pub run_id: i64, |
| 155 | } |
| 156 | |
| 157 | #[async_trait::async_trait] |
| 158 | impl ArtifactSink for UploadSink<'_> { |
| 159 | async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String> { |
| 160 | let url = self.client.url(&format!( |
| 161 | "/-/runner/jobs/{}/artifacts/{}", |
| 162 | self.run_id, spec.name |
| 163 | )); |
| 164 | let resp = self |
| 165 | .client |
| 166 | .auth( |
| 167 | self.client |
| 168 | .http |
| 169 | .post(url) |
| 170 | .header("Content-Type", "application/octet-stream") |
| 171 | .body(tar.to_vec()), |
| 172 | ) |
| 173 | .send() |
| 174 | .await |
| 175 | .map_err(|e| format!("upload: {e}"))?; |
| 176 | if !resp.status().is_success() { |
| 177 | return Err(format!( |
| 178 | "upload: {} {}", |
| 179 | resp.status(), |
| 180 | body_hint(resp).await |
| 181 | )); |
| 182 | } |
| 183 | resp.json::<Stored>() |
| 184 | .await |
| 185 | .map_err(|e| format!("decoding upload response: {e}")) |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | /// A short slice of an error response, for a log line. Bounded because the |
| 190 | /// body could be an HTML error page. |
| 191 | async fn body_hint(resp: reqwest::Response) -> String { |
| 192 | let text = resp.text().await.unwrap_or_default(); |
| 193 | let trimmed = text.trim(); |
| 194 | if trimmed.len() > 200 { |
| 195 | format!("{}…", &trimmed[..200]) |
| 196 | } else { |
| 197 | trimmed.to_string() |
| 198 | } |
| 199 | } |