anvilsign in

collin/anvil

main / crates / anvil-web / src / runner.rs
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
13use anvil_core::App;
14use anvil_job::{
15 JobResult,
16 RunnerInfo,
17};
18use 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.
44const UPLOAD_HARD_CAP: usize = 2 * 1024 * 1024 * 1024;
45
46pub 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.
61fn 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.
88fn 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.
96fn 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.
116async 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, &info.version);
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(anvil_job::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/// The checkout for a claimed job, as an uncompressed tar.
138async fn checkout(
139 State(app): State<App>,
140 Path(run_id): Path<i64>,
141 headers: HeaderMap,
142) -> Result<Response, Response> {
143 // A GET carries no body, so the runner name comes from the header pair
144 // rather than a JSON payload.
145 let info = info_from_headers(&headers);
146 authorize_job(&app, &headers, &info, run_id)?;
147 let tar = anvil_ci::checkout_tar(&app, run_id)
148 .await
149 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
150 Ok(([("content-type", "application/x-tar")], tar).into_response())
151}
152
153async fn heartbeat(
154 State(app): State<App>,
155 Path(run_id): Path<i64>,
156 headers: HeaderMap,
157 axum::Json(info): axum::Json<RunnerInfo>,
158) -> Result<Response, Response> {
159 let name = authenticate(&app, &headers, &info)?;
160 // A runner mid-build is not claiming, so its heartbeats are the only thing
161 // keeping it visible to platform routing while a long job runs.
162 app.jobs.seen(&name, &info.platform, &info.version);
163 if app.jobs.heartbeat(run_id, &name) {
164 Ok(StatusCode::NO_CONTENT.into_response())
165 } else {
166 Err((StatusCode::CONFLICT, "you do not hold this run").into_response())
167 }
168}
169
170/// Store one artifact tar, returning what actually landed on disk so the
171/// runner can charge its run budget.
172async fn artifact(
173 State(app): State<App>,
174 Path((run_id, name)): Path<(i64, String)>,
175 headers: HeaderMap,
176 body: Bytes,
177) -> Result<Response, Response> {
178 let info = info_from_headers(&headers);
179 authorize_job(&app, &headers, &info, run_id)?;
180
181 // The spec is re-derived from the commit's pipeline, never taken from the
182 // request: `browse` decides whether this tar is extracted into a servable
183 // directory tree, which is not the runner's call to make.
184 let spec = anvil_ci::artifact_spec(&app, run_id, &name)
185 .await
186 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?
187 .ok_or_else(|| {
188 (
189 StatusCode::BAD_REQUEST,
190 format!("{name} is not a declared artifact of this run"),
191 )
192 .into_response()
193 })?;
194
195 let run = anvil_core::ci::get(&app.db, run_id)
196 .await
197 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response())?
198 .ok_or_else(|| StatusCode::NOT_FOUND.into_response())?;
199
200 let stored = anvil_ci::store_upload(&app, &run, &spec, &body)
201 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
202 Ok(axum::Json(stored).into_response())
203}
204
205/// Record a finished job.
206async fn result(
207 State(app): State<App>,
208 Path(run_id): Path<i64>,
209 headers: HeaderMap,
210 axum::Json(result): axum::Json<JobResult>,
211) -> Result<Response, Response> {
212 let info = info_from_headers(&headers);
213 authorize_job(&app, &headers, &info, run_id)?;
214 anvil_ci::finish_run(&app, run_id, result)
215 .await
216 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
217 Ok(StatusCode::NO_CONTENT.into_response())
218}
219
220/// Runner identity for the endpoints with no JSON body to carry it.
221fn info_from_headers(headers: &HeaderMap) -> RunnerInfo {
222 RunnerInfo {
223 name: headers
224 .get("X-Anvil-Runner-Name")
225 .and_then(|v| v.to_str().ok())
226 .unwrap_or_default()
227 .to_string(),
228 platform: String::new(),
229 version: String::new(),
230 }
231}