| 1 | //! Endpoints a job runner dials in to (see `docs/remote-runners.md`). |
| 2 | //! |
| 3 | //! anvil does not execute CI; runners claim jobs here and report back. All |
| 4 | //! five endpoints authenticate with `[ci] runner_token` in the |
| 5 | //! `X-Anvil-Runner-Token` header, and the four per-job ones additionally |
| 6 | //! require the caller to hold that run's lease — a valid token gets you *a* |
| 7 | //! job, not everyone else's. |
| 8 | //! |
| 9 | //! Deliberately outside the browser session/CSRF world: these are machine |
| 10 | //! callers with a bearer secret, never a cookie, so none of the mutating |
| 11 | //! routes here take a CSRF token. |
| 12 | |
| 13 | use anvil_core::App; |
| 14 | use anvil_job::{ |
| 15 | JobResult, |
| 16 | RunnerInfo, |
| 17 | }; |
| 18 | use axum::{ |
| 19 | Router, |
| 20 | body::Bytes, |
| 21 | extract::{ |
| 22 | Path, |
| 23 | State, |
| 24 | }, |
| 25 | http::{ |
| 26 | HeaderMap, |
| 27 | StatusCode, |
| 28 | }, |
| 29 | response::{ |
| 30 | IntoResponse, |
| 31 | Response, |
| 32 | }, |
| 33 | routing::{ |
| 34 | get, |
| 35 | post, |
| 36 | }, |
| 37 | }; |
| 38 | |
| 39 | /// Cap on an uploaded artifact tar, over and above `ci.artifact_max_mb`. |
| 40 | /// |
| 41 | /// The runner already applies the configured cap while downloading, so this is |
| 42 | /// a backstop against a runner that ignores it, not the real limit. Generous |
| 43 | /// enough never to bite a legitimate upload. |
| 44 | const UPLOAD_HARD_CAP: usize = 2 * 1024 * 1024 * 1024; |
| 45 | |
| 46 | pub fn routes(router: Router<App>) -> Router<App> { |
| 47 | router |
| 48 | .route("/-/runner/claim", post(claim)) |
| 49 | .route("/-/runner/jobs/{run_id}/checkout.tar", get(checkout)) |
| 50 | .route("/-/runner/jobs/{run_id}/heartbeat", post(heartbeat)) |
| 51 | .route("/-/runner/jobs/{run_id}/artifacts/{name}", post(artifact)) |
| 52 | .route("/-/runner/jobs/{run_id}/result", post(result)) |
| 53 | .layer(axum::extract::DefaultBodyLimit::max(UPLOAD_HARD_CAP)) |
| 54 | } |
| 55 | |
| 56 | /// Authenticate a runner and return the name it claims. |
| 57 | /// |
| 58 | /// The name is self-asserted and is *not* a credential — it labels runs and |
| 59 | /// keys leases. Everyone holding the token is the same principal, which is the |
| 60 | /// documented shape of `runner_token` until per-runner credentials exist. |
| 61 | fn authenticate(app: &App, headers: &HeaderMap, info: &RunnerInfo) -> Result<String, Response> { |
| 62 | let configured = app.config.ci.runner_token.as_bytes(); |
| 63 | if configured.is_empty() { |
| 64 | // Not "forbidden": nothing the caller can fix. Say so plainly rather |
| 65 | // than leaving an operator guessing at a 403. |
| 66 | return Err(( |
| 67 | StatusCode::SERVICE_UNAVAILABLE, |
| 68 | "[ci] runner_token is not set on this instance", |
| 69 | ) |
| 70 | .into_response()); |
| 71 | } |
| 72 | let presented = headers |
| 73 | .get("X-Anvil-Runner-Token") |
| 74 | .map(|v| v.as_bytes()) |
| 75 | .unwrap_or_default(); |
| 76 | if !constant_time_eq(presented, configured) { |
| 77 | return Err(StatusCode::UNAUTHORIZED.into_response()); |
| 78 | } |
| 79 | let name = info.name.trim(); |
| 80 | if name.is_empty() { |
| 81 | return Err((StatusCode::BAD_REQUEST, "runner name is required").into_response()); |
| 82 | } |
| 83 | Ok(name.to_string()) |
| 84 | } |
| 85 | |
| 86 | /// Compare without leaking where the mismatch is, matching how CSRF tokens are |
| 87 | /// checked elsewhere in the crate. |
| 88 | fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { |
| 89 | if a.len() != b.len() { |
| 90 | return false; |
| 91 | } |
| 92 | a.iter().zip(b).fold(0u8, |acc, (x, y)| acc | (x ^ y)) == 0 |
| 93 | } |
| 94 | |
| 95 | /// Check the token *and* that this runner holds the run. |
| 96 | fn authorize_job( |
| 97 | app: &App, |
| 98 | headers: &HeaderMap, |
| 99 | info: &RunnerInfo, |
| 100 | run_id: i64, |
| 101 | ) -> Result<String, Response> { |
| 102 | let name = authenticate(app, headers, info)?; |
| 103 | if !app.jobs.holds(run_id, &name) { |
| 104 | // 409 rather than 403: the caller is a legitimate runner, the lease |
| 105 | // just is not theirs any more (expired and requeued, most likely). |
| 106 | return Err((StatusCode::CONFLICT, "you do not hold this run").into_response()); |
| 107 | } |
| 108 | Ok(name) |
| 109 | } |
| 110 | |
| 111 | /// Long-poll for a job. 204 when the poll expires with nothing queued. |
| 112 | /// |
| 113 | /// Parking here rather than returning immediately is what keeps an idle runner |
| 114 | /// from hammering the forge: one request per minute at rest, and a job starts |
| 115 | /// the moment a push enqueues one. |
| 116 | async fn claim( |
| 117 | State(app): State<App>, |
| 118 | headers: HeaderMap, |
| 119 | axum::Json(info): axum::Json<RunnerInfo>, |
| 120 | ) -> Result<Response, Response> { |
| 121 | let name = authenticate(&app, &headers, &info)?; |
| 122 | // Note the runner *before* claiming: what platforms are present decides |
| 123 | // which jobs are routed where, and a runner asking for work is the freshest |
| 124 | // evidence there is that its architecture has someone behind it. |
| 125 | app.jobs.seen(&name, &info.platform); |
| 126 | |
| 127 | if let Some(job) = anvil_ci::claim_next(&app, &name, &info.platform).await { |
| 128 | return Ok(axum::Json(job).into_response()); |
| 129 | } |
| 130 | app.jobs.wait_for_work(CLAIM_POLL).await; |
| 131 | match anvil_ci::claim_next(&app, &name, &info.platform).await { |
| 132 | Some(job) => Ok(axum::Json(job).into_response()), |
| 133 | None => Ok(StatusCode::NO_CONTENT.into_response()), |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | /// Matches the runner's own poll window; it gives up a shade later. |
| 138 | const CLAIM_POLL: std::time::Duration = std::time::Duration::from_secs(55); |
| 139 | |
| 140 | /// The checkout for a claimed job, as an uncompressed tar. |
| 141 | async fn checkout( |
| 142 | State(app): State<App>, |
| 143 | Path(run_id): Path<i64>, |
| 144 | headers: HeaderMap, |
| 145 | ) -> Result<Response, Response> { |
| 146 | // A GET carries no body, so the runner name comes from the header pair |
| 147 | // rather than a JSON payload. |
| 148 | let info = info_from_headers(&headers); |
| 149 | authorize_job(&app, &headers, &info, run_id)?; |
| 150 | let tar = anvil_ci::checkout_tar(&app, run_id) |
| 151 | .await |
| 152 | .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?; |
| 153 | Ok(([("content-type", "application/x-tar")], tar).into_response()) |
| 154 | } |
| 155 | |
| 156 | async fn heartbeat( |
| 157 | State(app): State<App>, |
| 158 | Path(run_id): Path<i64>, |
| 159 | headers: HeaderMap, |
| 160 | axum::Json(info): axum::Json<RunnerInfo>, |
| 161 | ) -> Result<Response, Response> { |
| 162 | let name = authenticate(&app, &headers, &info)?; |
| 163 | // A runner mid-build is not claiming, so its heartbeats are the only thing |
| 164 | // keeping it visible to platform routing while a long job runs. |
| 165 | app.jobs.seen(&name, &info.platform); |
| 166 | if app.jobs.heartbeat(run_id, &name) { |
| 167 | Ok(StatusCode::NO_CONTENT.into_response()) |
| 168 | } else { |
| 169 | Err((StatusCode::CONFLICT, "you do not hold this run").into_response()) |
| 170 | } |
| 171 | } |
| 172 | |
| 173 | /// Store one artifact tar, returning what actually landed on disk so the |
| 174 | /// runner can charge its run budget. |
| 175 | async fn artifact( |
| 176 | State(app): State<App>, |
| 177 | Path((run_id, name)): Path<(i64, String)>, |
| 178 | headers: HeaderMap, |
| 179 | body: Bytes, |
| 180 | ) -> Result<Response, Response> { |
| 181 | let info = info_from_headers(&headers); |
| 182 | authorize_job(&app, &headers, &info, run_id)?; |
| 183 | |
| 184 | // The spec is re-derived from the commit's pipeline, never taken from the |
| 185 | // request: `browse` decides whether this tar is extracted into a servable |
| 186 | // directory tree, which is not the runner's call to make. |
| 187 | let spec = anvil_ci::artifact_spec(&app, run_id, &name) |
| 188 | .await |
| 189 | .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())? |
| 190 | .ok_or_else(|| { |
| 191 | ( |
| 192 | StatusCode::BAD_REQUEST, |
| 193 | format!("{name} is not a declared artifact of this run"), |
| 194 | ) |
| 195 | .into_response() |
| 196 | })?; |
| 197 | |
| 198 | let run = anvil_core::ci::get(&app.db, run_id) |
| 199 | .await |
| 200 | .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response())? |
| 201 | .ok_or_else(|| StatusCode::NOT_FOUND.into_response())?; |
| 202 | |
| 203 | let stored = anvil_ci::store_upload(&app, &run, &spec, &body) |
| 204 | .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?; |
| 205 | Ok(axum::Json(stored).into_response()) |
| 206 | } |
| 207 | |
| 208 | /// Record a finished job. |
| 209 | async fn result( |
| 210 | State(app): State<App>, |
| 211 | Path(run_id): Path<i64>, |
| 212 | headers: HeaderMap, |
| 213 | axum::Json(result): axum::Json<JobResult>, |
| 214 | ) -> Result<Response, Response> { |
| 215 | let info = info_from_headers(&headers); |
| 216 | authorize_job(&app, &headers, &info, run_id)?; |
| 217 | anvil_ci::finish_run(&app, run_id, result) |
| 218 | .await |
| 219 | .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?; |
| 220 | Ok(StatusCode::NO_CONTENT.into_response()) |
| 221 | } |
| 222 | |
| 223 | /// Runner identity for the endpoints with no JSON body to carry it. |
| 224 | fn info_from_headers(headers: &HeaderMap) -> RunnerInfo { |
| 225 | RunnerInfo { |
| 226 | name: headers |
| 227 | .get("X-Anvil-Runner-Name") |
| 228 | .and_then(|v| v.to_str().ok()) |
| 229 | .unwrap_or_default() |
| 230 | .to_string(), |
| 231 | platform: String::new(), |
| 232 | version: String::new(), |
| 233 | } |
| 234 | } |