| 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 | pub fn new(base: &str, token: &str, info: RunnerInfo) -> Result<Self, String> { |
| 33 | let http = reqwest::Client::builder() |
| 34 | .timeout(REQUEST_TIMEOUT) |
| 35 | .build() |
| 36 | .map_err(|e| format!("building http client: {e}"))?; |
| 37 | Ok(Self { |
| 38 | http, |
| 39 | base: base.trim_end_matches('/').to_string(), |
| 40 | token: token.to_string(), |
| 41 | info, |
| 42 | }) |
| 43 | } |
| 44 | |
| 45 | fn url(&self, path: &str) -> String { |
| 46 | format!("{}{path}", self.base) |
| 47 | } |
| 48 | |
| 49 | /// Token plus identity on every request. The name goes in a header rather |
| 50 | /// than only in the claim body so the endpoints with no body of their own |
| 51 | /// (the checkout GET, artifact uploads) can be authorized the same way. |
| 52 | fn auth(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder { |
| 53 | req.header("X-Anvil-Runner-Token", &self.token) |
| 54 | .header("X-Anvil-Runner-Name", &self.info.name) |
| 55 | } |
| 56 | |
| 57 | /// Long-poll for a job. `Ok(None)` means the poll expired with nothing |
| 58 | /// queued, which is the common case and not an error. |
| 59 | pub async fn claim(&self) -> Result<Option<JobSpec>, String> { |
| 60 | let resp = self |
| 61 | .auth(self.http.post(self.url("/-/runner/claim"))) |
| 62 | .json(&self.info) |
| 63 | .send() |
| 64 | .await |
| 65 | .map_err(|e| format!("claim: {e}"))?; |
| 66 | match resp.status() { |
| 67 | reqwest::StatusCode::NO_CONTENT => Ok(None), |
| 68 | s if s.is_success() => resp |
| 69 | .json::<JobSpec>() |
| 70 | .await |
| 71 | .map(Some) |
| 72 | .map_err(|e| format!("decoding job: {e}")), |
| 73 | s => Err(format!("claim: {s} {}", body_hint(resp).await)), |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | pub async fn checkout(&self, run_id: i64) -> Result<Vec<u8>, String> { |
| 78 | let resp = self |
| 79 | .auth( |
| 80 | self.http |
| 81 | .get(self.url(&format!("/-/runner/jobs/{run_id}/checkout.tar"))), |
| 82 | ) |
| 83 | .send() |
| 84 | .await |
| 85 | .map_err(|e| format!("checkout: {e}"))?; |
| 86 | if !resp.status().is_success() { |
| 87 | return Err(format!( |
| 88 | "checkout: {} {}", |
| 89 | resp.status(), |
| 90 | body_hint(resp).await |
| 91 | )); |
| 92 | } |
| 93 | resp.bytes() |
| 94 | .await |
| 95 | .map(|b| b.to_vec()) |
| 96 | .map_err(|e| format!("reading checkout: {e}")) |
| 97 | } |
| 98 | |
| 99 | /// Extend the lease. `false` means anvil no longer thinks we hold this run |
| 100 | /// — it was requeued out from under us, and the job should be abandoned |
| 101 | /// rather than reported. |
| 102 | pub async fn heartbeat(&self, run_id: i64) -> bool { |
| 103 | let sent = self |
| 104 | .auth( |
| 105 | self.http |
| 106 | .post(self.url(&format!("/-/runner/jobs/{run_id}/heartbeat"))) |
| 107 | .json(&self.info), |
| 108 | ) |
| 109 | .send() |
| 110 | .await; |
| 111 | matches!(sent, Ok(r) if r.status().is_success()) |
| 112 | } |
| 113 | |
| 114 | pub async fn report(&self, run_id: i64, result: &JobResult) -> Result<(), String> { |
| 115 | let resp = self |
| 116 | .auth( |
| 117 | self.http |
| 118 | .post(self.url(&format!("/-/runner/jobs/{run_id}/result"))) |
| 119 | .json(result), |
| 120 | ) |
| 121 | .send() |
| 122 | .await |
| 123 | .map_err(|e| format!("result: {e}"))?; |
| 124 | if resp.status().is_success() { |
| 125 | Ok(()) |
| 126 | } else { |
| 127 | Err(format!( |
| 128 | "result: {} {}", |
| 129 | resp.status(), |
| 130 | body_hint(resp).await |
| 131 | )) |
| 132 | } |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | /// Uploads each artifact tar to anvil as it comes out of the container. |
| 137 | /// |
| 138 | /// The server does the storing, and reports back what it wrote — the runner |
| 139 | /// never learns anvil's on-disk layout, and `browse` stays the server's call. |
| 140 | pub struct UploadSink<'a> { |
| 141 | pub client: &'a Client, |
| 142 | pub run_id: i64, |
| 143 | } |
| 144 | |
| 145 | #[async_trait::async_trait] |
| 146 | impl ArtifactSink for UploadSink<'_> { |
| 147 | async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String> { |
| 148 | let url = self.client.url(&format!( |
| 149 | "/-/runner/jobs/{}/artifacts/{}", |
| 150 | self.run_id, spec.name |
| 151 | )); |
| 152 | let resp = self |
| 153 | .client |
| 154 | .auth( |
| 155 | self.client |
| 156 | .http |
| 157 | .post(url) |
| 158 | .header("Content-Type", "application/octet-stream") |
| 159 | .body(tar.to_vec()), |
| 160 | ) |
| 161 | .send() |
| 162 | .await |
| 163 | .map_err(|e| format!("upload: {e}"))?; |
| 164 | if !resp.status().is_success() { |
| 165 | return Err(format!( |
| 166 | "upload: {} {}", |
| 167 | resp.status(), |
| 168 | body_hint(resp).await |
| 169 | )); |
| 170 | } |
| 171 | resp.json::<Stored>() |
| 172 | .await |
| 173 | .map_err(|e| format!("decoding upload response: {e}")) |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | /// A short slice of an error response, for a log line. Bounded because the |
| 178 | /// body could be an HTML error page. |
| 179 | async fn body_hint(resp: reqwest::Response) -> String { |
| 180 | let text = resp.text().await.unwrap_or_default(); |
| 181 | let trimmed = text.trim(); |
| 182 | if trimmed.len() > 200 { |
| 183 | format!("{}…", &trimmed[..200]) |
| 184 | } else { |
| 185 | trimmed.to_string() |
| 186 | } |
| 187 | } |