anvilsign in

collin/anvil

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);
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.
138const CLAIM_POLL: std::time::Duration = std::time::Duration::from_secs(55);
139
140/// The checkout for a claimed job, as an uncompressed tar.
141async 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
156async 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.
175async 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.
209async 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.
224fn 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}