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
123 if let Some(job) = anvil_ci::claim_next(&app, &name).await {
124 return Ok(axum::Json(job).into_response());
125 }
126 app.jobs.wait_for_work(CLAIM_POLL).await;
127 match anvil_ci::claim_next(&app, &name).await {
128 Some(job) => Ok(axum::Json(job).into_response()),
129 None => Ok(StatusCode::NO_CONTENT.into_response()),
130 }
131}
132
133/// Matches the runner's own poll window; it gives up a shade later.
134const CLAIM_POLL: std::time::Duration = std::time::Duration::from_secs(55);
135
136/// The checkout for a claimed job, as an uncompressed tar.
137async fn checkout(
138 State(app): State<App>,
139 Path(run_id): Path<i64>,
140 headers: HeaderMap,
141) -> Result<Response, Response> {
142 // A GET carries no body, so the runner name comes from the header pair
143 // rather than a JSON payload.
144 let info = info_from_headers(&headers);
145 authorize_job(&app, &headers, &info, run_id)?;
146 let tar = anvil_ci::checkout_tar(&app, run_id)
147 .await
148 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
149 Ok(([("content-type", "application/x-tar")], tar).into_response())
150}
151
152async fn heartbeat(
153 State(app): State<App>,
154 Path(run_id): Path<i64>,
155 headers: HeaderMap,
156 axum::Json(info): axum::Json<RunnerInfo>,
157) -> Result<Response, Response> {
158 let name = authenticate(&app, &headers, &info)?;
159 if app.jobs.heartbeat(run_id, &name) {
160 Ok(StatusCode::NO_CONTENT.into_response())
161 } else {
162 Err((StatusCode::CONFLICT, "you do not hold this run").into_response())
163 }
164}
165
166/// Store one artifact tar, returning what actually landed on disk so the
167/// runner can charge its run budget.
168async fn artifact(
169 State(app): State<App>,
170 Path((run_id, name)): Path<(i64, String)>,
171 headers: HeaderMap,
172 body: Bytes,
173) -> Result<Response, Response> {
174 let info = info_from_headers(&headers);
175 authorize_job(&app, &headers, &info, run_id)?;
176
177 // The spec is re-derived from the commit's pipeline, never taken from the
178 // request: `browse` decides whether this tar is extracted into a servable
179 // directory tree, which is not the runner's call to make.
180 let spec = anvil_ci::artifact_spec(&app, run_id, &name)
181 .await
182 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?
183 .ok_or_else(|| {
184 (
185 StatusCode::BAD_REQUEST,
186 format!("{name} is not a declared artifact of this run"),
187 )
188 .into_response()
189 })?;
190
191 let run = anvil_core::ci::get(&app.db, run_id)
192 .await
193 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response())?
194 .ok_or_else(|| StatusCode::NOT_FOUND.into_response())?;
195
196 let stored = anvil_ci::store_upload(&app, &run, &spec, &body)
197 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
198 Ok(axum::Json(stored).into_response())
199}
200
201/// Record a finished job.
202async fn result(
203 State(app): State<App>,
204 Path(run_id): Path<i64>,
205 headers: HeaderMap,
206 axum::Json(result): axum::Json<JobResult>,
207) -> Result<Response, Response> {
208 let info = info_from_headers(&headers);
209 authorize_job(&app, &headers, &info, run_id)?;
210 anvil_ci::finish_run(&app, run_id, result)
211 .await
212 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e).into_response())?;
213 Ok(StatusCode::NO_CONTENT.into_response())
214}
215
216/// Runner identity for the endpoints with no JSON body to carry it.
217fn info_from_headers(headers: &HeaderMap) -> RunnerInfo {
218 RunnerInfo {
219 name: headers
220 .get("X-Anvil-Runner-Name")
221 .and_then(|v| v.to_str().ok())
222 .unwrap_or_default()
223 .to_string(),
224 platform: String::new(),
225 version: String::new(),
226 }
227}