collin/anvil · 9d6f3bcd
Move CI execution into a dial-out anvil-worker runner
Collin Richards · 2026-08-24 08:03 UTC · 9d6f3bcda3e37c93b2eb302a1bb8371d16b082fd · parent b3adba61 · browse files
modifiedCargo.lock+21 −3
| ⋯ 141 unchanged lines | |||
| 142 | 142 | "anvil-docker", | |
| 143 | 143 | "anvil-git", | |
| 144 | 144 | "anvil-job", | |
| 145 | - | "async-trait", | |
| 146 | - | "bollard", | |
| 147 | 145 | "flate2", | |
| 148 | - | "futures-util", | |
| 149 | 146 | "reqwest", | |
| 150 | 147 | "serde_json", | |
| 151 | 148 | "tar", | |
| ⋯ 30 unchanged lines | |||
| 182 | 179 | version = "0.0.0" | |
| 183 | 180 | dependencies = [ | |
| 184 | 181 | "aes-gcm", | |
| 182 | + | "anvil-job", | |
| 185 | 183 | "argon2 0.5.3", | |
| 186 | 184 | "async-trait", | |
| 187 | 185 | "base64", | |
| ⋯ 67 unchanged lines | |||
| 255 | 253 | version = "0.0.0" | |
| 256 | 254 | dependencies = [ | |
| 257 | 255 | "anvil-agent", | |
| 256 | + | "anvil-ci", | |
| 258 | 257 | "anvil-core", | |
| 259 | 258 | "anvil-git", | |
| 259 | + | "anvil-job", | |
| 260 | 260 | "argon2 0.5.3", | |
| 261 | 261 | "axum", | |
| 262 | 262 | "axum-extra", | |
| ⋯ 22 unchanged lines | |||
| 285 | 285 | ] | |
| 286 | 286 | ||
| 287 | 287 | [[package]] | |
| 288 | + | name = "anvil-worker" | |
| 289 | + | version = "0.0.0" | |
| 290 | + | dependencies = [ | |
| 291 | + | "anvil-docker", | |
| 292 | + | "anvil-job", | |
| 293 | + | "async-trait", | |
| 294 | + | "bollard", | |
| 295 | + | "clap", | |
| 296 | + | "futures-util", | |
| 297 | + | "reqwest", | |
| 298 | + | "serde_json", | |
| 299 | + | "tar", | |
| 300 | + | "tokio", | |
| 301 | + | "tracing", | |
| 302 | + | "tracing-subscriber", | |
| 303 | + | ] | |
| 304 | + | ||
| 305 | + | [[package]] | |
| 288 | 306 | name = "anyhow" | |
| 289 | 307 | version = "1.0.102" | |
| 290 | 308 | source = "registry+https://github.com/rust-lang/crates.io-index" | |
| ⋯ 5510 unchanged lines | |||
modifiedcrates/anvil-ci/Cargo.toml+0 −3
| ⋯ 11 unchanged lines | |||
| 12 | 12 | anvil-docker.workspace = true | |
| 13 | 13 | anvil-job.workspace = true | |
| 14 | 14 | anvil-git.workspace = true | |
| 15 | - | async-trait.workspace = true | |
| 16 | - | bollard.workspace = true | |
| 17 | 15 | flate2.workspace = true | |
| 18 | - | futures-util.workspace = true | |
| 19 | 16 | reqwest.workspace = true | |
| 20 | 17 | serde_json.workspace = true | |
| 21 | 18 | tar.workspace = true | |
| ⋯ 5 unchanged lines | |||
modifiedcrates/anvil-ci/src/lib.rs+229 −399
| 1 | - | //! CI runner: drains queued [`anvil_core::ci`] runs and executes their | |
| 2 | - | //! `.anvil/ci.yml` pipeline in a Docker container (via the socket, using | |
| 3 | - | //! bollard). | |
| 1 | + | //! CI dispatch: decides what a queued [`anvil_core::ci`] run should do, hands | |
| 2 | + | //! it to a runner that dialled in, and records what came back. | |
| 4 | 3 | //! | |
| 5 | - | //! For each run: resolve the repo, materialize the commit's tree, parse the | |
| 6 | - | //! pipeline, then run all steps as one `set -e` shell script inside the | |
| 7 | - | //! pipeline's image. The checkout is uploaded into the container as a tar (via | |
| 8 | - | //! the Docker API), so it works regardless of where anvil's own filesystem | |
| 9 | - | //! lives and never exposes anvil's data volume to CI. | |
| 4 | + | //! anvil does **not** execute jobs. It used to — this crate was the runner — | |
| 5 | + | //! but a release build needs more RAM than the host it runs on has, so | |
| 6 | + | //! execution moved to `anvil-worker` on a machine that can afford it. See | |
| 7 | + | //! `docs/remote-runners.md`. | |
| 10 | 8 | //! | |
| 11 | - | //! anvil is the *broker*: it is the only Docker client, and the job container | |
| 12 | - | //! gets no socket, no bind mounts, and no volumes. On top of that the job runs | |
| 13 | - | //! with all capabilities dropped, `no-new-privileges`, and configurable | |
| 14 | - | //! pids/memory/cpu caps plus a wall-clock timeout and an optional image | |
| 15 | - | //! allowlist ([`anvil_core::config::CiConfig`]). | |
| 9 | + | //! What stays here is everything a runner must not be trusted with: resolving | |
| 10 | + | //! the repo, materializing the commit's tree, parsing `.anvil/ci.yml`, checking | |
| 11 | + | //! the image allowlist, opening the secret vault, assembling the shell script, | |
| 12 | + | //! and deciding where artifacts land on disk. A runner receives an image, a | |
| 13 | + | //! script, a tar and some limits ([`anvil_job::JobSpec`]) and can neither widen | |
| 14 | + | //! its sandbox nor reinterpret a pipeline. | |
| 16 | 15 | ||
| 17 | 16 | use std::{ | |
| 18 | 17 | collections::BTreeMap, | |
| ⋯ 18 unchanged lines | |||
| 37 | 36 | }; | |
| 38 | 37 | use anvil_job::{ | |
| 39 | 38 | ArtifactSpec, | |
| 40 | - | CollectedArtifact, | |
| 39 | + | JobResult, | |
| 41 | 40 | JobSpec, | |
| 42 | 41 | META_DIR, | |
| 43 | - | META_TAR_CAP, | |
| 44 | - | META_VALUE_CAP, | |
| 45 | 42 | Sandbox, | |
| 46 | 43 | Stored, | |
| 47 | - | WORKDIR, | |
| 44 | + | mb_cap, | |
| 48 | 45 | }; | |
| 49 | - | use bollard::{ | |
| 50 | - | Docker, | |
| 51 | - | container::{ | |
| 52 | - | Config, | |
| 53 | - | CreateContainerOptions, | |
| 54 | - | DownloadFromContainerOptions, | |
| 55 | - | LogsOptions, | |
| 56 | - | RemoveContainerOptions, | |
| 57 | - | StartContainerOptions, | |
| 58 | - | UploadToContainerOptions, | |
| 59 | - | WaitContainerOptions, | |
| 60 | - | }, | |
| 61 | - | models::HostConfig, | |
| 62 | - | }; | |
| 63 | - | use futures_util::StreamExt; | |
| 64 | 46 | use tokio::sync::mpsc::UnboundedReceiver; | |
| 65 | 47 | ||
| 66 | 48 | /// Turn a parsed pipeline into the job a runner receives. | |
| ⋯ 83 unchanged lines | |||
| 150 | 132 | script | |
| 151 | 133 | } | |
| 152 | 134 | ||
| 153 | - | /// Run the CI worker loop: recover interrupted runs, drain the queue, then | |
| 154 | - | /// process run ids as they arrive on `rx`. Runs one job at a time. | |
| 155 | - | pub async fn run_worker(app: App, mut rx: UnboundedReceiver<i64>) { | |
| 135 | + | /// Dispatch queued runs to runners: recover interrupted runs, then wake parked | |
| 136 | + | /// claim requests as work arrives, and requeue jobs whose runner went away. | |
| 137 | + | /// | |
| 138 | + | /// anvild no longer executes anything (see `docs/remote-runners.md`). This | |
| 139 | + | /// loop used to *be* the runner; now it only decides that a run is ready to be | |
| 140 | + | /// claimed. `rx` still carries newly-enqueued run ids from every push path, so | |
| 141 | + | /// nothing that enqueues had to change. | |
| 142 | + | pub async fn run_dispatcher(app: App, mut rx: UnboundedReceiver<i64>) { | |
| 156 | 143 | match ci::requeue_interrupted(&app.db).await { | |
| 157 | 144 | Ok(ids) if !ids.is_empty() => { | |
| 158 | 145 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) | |
| ⋯ 1 unchanged line | |||
| 160 | 147 | Ok(_) => {} | |
| 161 | 148 | Err(e) => tracing::error!("ci: requeue failed: {e}"), | |
| 162 | 149 | } | |
| 163 | - | match ci::queued_ids(&app.db).await { | |
| 164 | - | Ok(ids) => { | |
| 165 | - | for id in ids { | |
| 166 | - | run_one(&app, id).await; | |
| 150 | + | ||
| 151 | + | tokio::spawn(sweep_expired_leases(app.clone())); | |
| 152 | + | ||
| 153 | + | if app.config.ci.runner_token.is_empty() { | |
| 154 | + | tracing::warn!( | |
| 155 | + | "ci: [ci] runner_token is empty, so no runner can authenticate — \ | |
| 156 | + | queued runs will sit until one is set" | |
| 157 | + | ); | |
| 158 | + | } | |
| 159 | + | tracing::info!("ci dispatcher ready"); | |
| 160 | + | ||
| 161 | + | // Anything already queued is claimable the moment a runner asks; the wake | |
| 162 | + | // is only for runners already parked in a long poll. | |
| 163 | + | app.jobs.wake(); | |
| 164 | + | while rx.recv().await.is_some() { | |
| 165 | + | app.jobs.wake(); | |
| 166 | + | } | |
| 167 | + | } | |
| 168 | + | ||
| 169 | + | /// Requeue runs whose runner stopped heartbeating. | |
| 170 | + | /// | |
| 171 | + | /// The case anvild's own crash never had: with execution in-process, a job | |
| 172 | + | /// died exactly when the process did, and `requeue_interrupted` at startup | |
| 173 | + | /// covered it. A remote runner can vanish while anvil stays up. | |
| 174 | + | /// | |
| 175 | + | /// Note this can double-run a job whose runner is alive but unreachable — the | |
| 176 | + | /// container keeps going, and the requeued run may be claimed elsewhere. CI | |
| 177 | + | /// steps are assumed idempotent, and one job at a time makes it unlikely, but | |
| 178 | + | /// it is the honest failure mode of a lease without fencing. | |
| 179 | + | async fn sweep_expired_leases(app: App) { | |
| 180 | + | let mut tick = tokio::time::interval(anvil_core::jobs::LEASE_TTL / 4); | |
| 181 | + | loop { | |
| 182 | + | tick.tick().await; | |
| 183 | + | for (run_id, runner) in app.jobs.take_expired() { | |
| 184 | + | tracing::warn!("ci: run {run_id} lease expired (runner {runner}); requeueing"); | |
| 185 | + | if let Err(e) = ci::requeue(&app.db, run_id).await { | |
| 186 | + | tracing::error!("ci: requeueing {run_id} failed: {e}"); | |
| 187 | + | continue; | |
| 167 | 188 | } | |
| 189 | + | app.jobs.wake(); | |
| 168 | 190 | } | |
| 169 | - | Err(e) => tracing::error!("ci: listing queued runs failed: {e}"), | |
| 170 | - | } | |
| 171 | - | tracing::info!("ci runner ready"); | |
| 172 | - | while let Some(id) = rx.recv().await { | |
| 173 | - | run_one(&app, id).await; | |
| 174 | 191 | } | |
| 175 | 192 | } | |
| 176 | 193 | ||
| 177 | - | async fn run_one(app: &App, run_id: i64) { | |
| 178 | - | tracing::info!("ci: run {run_id} starting"); | |
| 179 | - | if let Err(e) = process(app, run_id).await { | |
| 180 | - | tracing::error!("ci: run {run_id} errored: {e}"); | |
| 181 | - | let _ = ci::append_log(&app.db, run_id, &format!("\n[runner error] {e}\n")).await; | |
| 182 | - | let _ = ci::finish(&app.db, run_id, ci::status::ERROR).await; | |
| 194 | + | /// Hand the oldest queued run to `runner`, or `None` when there is nothing to | |
| 195 | + | /// do. | |
| 196 | + | /// | |
| 197 | + | /// Runs that cannot be dispatched at all — an unparseable pipeline, a sealed | |
| 198 | + | /// vault, an image the allowlist forbids — are finished here rather than | |
| 199 | + | /// handed out, and the next queued run is tried. That keeps a single bad | |
| 200 | + | /// pipeline from wedging the queue. | |
| 201 | + | pub async fn claim_next(app: &App, runner: &str) -> Option<JobSpec> { | |
| 202 | + | let queued = match ci::queued_ids(&app.db).await { | |
| 203 | + | Ok(ids) => ids, | |
| 204 | + | Err(e) => { | |
| 205 | + | tracing::error!("ci: listing queued runs failed: {e}"); | |
| 206 | + | return None; | |
| 207 | + | } | |
| 208 | + | }; | |
| 209 | + | for run_id in queued { | |
| 210 | + | // Skip anything a concurrent claim already took: `queued_ids` reads | |
| 211 | + | // committed state, and `prepare` marks running. | |
| 212 | + | if app.jobs.holds_any(run_id) { | |
| 213 | + | continue; | |
| 214 | + | } | |
| 215 | + | match prepare(app, run_id, runner).await { | |
| 216 | + | Ok(Some(job)) => return Some(job), | |
| 217 | + | Ok(None) => continue, | |
| 218 | + | Err(e) => { | |
| 219 | + | tracing::error!("ci: preparing run {run_id} failed: {e}"); | |
| 220 | + | fail_run(app, run_id, &format!("\n[dispatch error] {e}\n")).await; | |
| 221 | + | } | |
| 222 | + | } | |
| 183 | 223 | } | |
| 224 | + | None | |
| 184 | 225 | } | |
| 185 | 226 | ||
| 186 | - | async fn process(app: &App, run_id: i64) -> Result<(), String> { | |
| 227 | + | /// Everything anvil does before a job leaves the building: resolve the repo, | |
| 228 | + | /// parse the pipeline, open the vault, build the spec, take the lease. | |
| 229 | + | /// | |
| 230 | + | /// `Ok(None)` means the run was finished here and should not be dispatched. | |
| 231 | + | async fn prepare(app: &App, run_id: i64, runner: &str) -> Result<Option<JobSpec>, String> { | |
| 187 | 232 | let run = ci::get(&app.db, run_id) | |
| 188 | 233 | .await | |
| 189 | 234 | .map_err(|e| e.to_string())? | |
| ⋯ 13 unchanged lines | |||
| 203 | 248 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; | |
| 204 | 249 | let pipeline = | |
| 205 | 250 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; | |
| 206 | - | ||
| 207 | - | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; | |
| 208 | - | let tar = build_tar(&files); | |
| 209 | - | ||
| 210 | - | ci::mark_running(&app.db, run_id).await.ok(); | |
| 211 | 251 | ||
| 212 | 252 | let short = &run.commit[..run.commit.len().min(12)]; | |
| 213 | 253 | let mut log = format!( | |
| 214 | - | "anvil ci · {}/{} · {} @ {short}\nimage: {}\n", | |
| 254 | + | "anvil ci · {}/{} · {} @ {short}\nrunner: {runner}\nimage: {}\n", | |
| 215 | 255 | owner.username, | |
| 216 | 256 | repo.name, | |
| 217 | 257 | run.ref_name, | |
| 218 | 258 | app.config.ci.resolve_image(&pipeline.image) | |
| 219 | 259 | ); | |
| 220 | 260 | ||
| 221 | - | // Artifacts are collected into a scratch directory next to their final | |
| 222 | - | // home (same filesystem, so the swap below is a rename), then moved into | |
| 223 | - | // place only after the run finishes. | |
| 224 | - | let artifacts_root = app.config.artifacts_dir(); | |
| 225 | - | let scratch = artifacts_root | |
| 226 | - | .join(run.repo_id.to_string()) | |
| 227 | - | .join(format!(".collecting-{run_id}")); | |
| 228 | - | ||
| 229 | 261 | // Secrets the pipeline asked for, from the in-memory vault. anvil holds no | |
| 230 | 262 | // key that opens the stored envelopes, so an unlock must have happened | |
| 231 | 263 | // (`anvild secret unlock`) or the run cannot proceed. | |
| ⋯ 11 unchanged lines | |||
| 243 | 275 | ci::append_log(&app.db, run_id, &log).await.ok(); | |
| 244 | 276 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); | |
| 245 | 277 | tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name); | |
| 246 | - | return Ok(()); | |
| 278 | + | return Ok(None); | |
| 247 | 279 | } | |
| 248 | 280 | }; | |
| 249 | 281 | if !env.is_empty() { | |
| ⋯ 6 unchanged lines | |||
| 256 | 288 | )); | |
| 257 | 289 | } | |
| 258 | 290 | ||
| 259 | - | let mut sink = ScratchSink { scratch: &scratch }; | |
| 260 | - | let outcome = match build_job(run_id, &pipeline, &app.config.ci, &env) { | |
| 261 | - | Ok(job) => execute(&job, tar, &mut log, &mut sink).await, | |
| 262 | - | Err(e) => Err(e), | |
| 291 | + | let job = match build_job(run_id, &pipeline, &app.config.ci, &env) { | |
| 292 | + | Ok(job) => job, | |
| 293 | + | Err(e) => { | |
| 294 | + | log.push_str(&format!("\n[runner error] {e}\n")); | |
| 295 | + | ci::append_log(&app.db, run_id, &log).await.ok(); | |
| 296 | + | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); | |
| 297 | + | return Ok(None); | |
| 298 | + | } | |
| 263 | 299 | }; | |
| 264 | - | let (status, collected) = match outcome { | |
| 265 | - | Ok((0, collected)) => (ci::status::SUCCESS, collected), | |
| 266 | - | Ok((code, collected)) => { | |
| 267 | - | log.push_str(&format!("\n[exited with status {code}]\n")); | |
| 268 | - | (ci::status::FAILURE, collected) | |
| 269 | - | } | |
| 270 | - | Err(e) => { | |
| 300 | + | ||
| 301 | + | // The header is written now rather than kept in memory until the result | |
| 302 | + | // lands, so a run that is visibly `running` has something to show. | |
| 303 | + | ci::append_log(&app.db, run_id, &log).await.ok(); | |
| 304 | + | ci::mark_running(&app.db, run_id).await.ok(); | |
| 305 | + | app.jobs.claim(run_id, runner, env); | |
| 306 | + | tracing::info!("ci: run {run_id} claimed by {runner}"); | |
| 307 | + | Ok(Some(job)) | |
| 308 | + | } | |
| 309 | + | ||
| 310 | + | /// The checkout a claimed job runs against, as an uncompressed tar. | |
| 311 | + | /// | |
| 312 | + | /// Rebuilt from the commit on demand rather than stashed at claim time: it is | |
| 313 | + | /// a pure function of the commit, so a runner retrying the fetch is free, and | |
| 314 | + | /// anvil holds no per-job buffer on a host where memory is the scarce thing. | |
| 315 | + | pub async fn checkout_tar(app: &App, run_id: i64) -> Result<Vec<u8>, String> { | |
| 316 | + | let (_, _, repo_path, run) = resolve(app, run_id).await?; | |
| 317 | + | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; | |
| 318 | + | Ok(build_tar(&files)) | |
| 319 | + | } | |
| 320 | + | ||
| 321 | + | /// The declared spec for one artifact of a run. | |
| 322 | + | /// | |
| 323 | + | /// Re-derived from the pipeline rather than trusted from the upload request: | |
| 324 | + | /// `browse` decides whether a tarball is extracted into a servable directory | |
| 325 | + | /// tree, which is not a choice a runner should get to make. | |
| 326 | + | pub async fn artifact_spec( | |
| 327 | + | app: &App, | |
| 328 | + | run_id: i64, | |
| 329 | + | name: &str, | |
| 330 | + | ) -> Result<Option<ArtifactSpec>, String> { | |
| 331 | + | let (_, _, repo_path, run) = resolve(app, run_id).await?; | |
| 332 | + | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) | |
| 333 | + | .map_err(|e| e.to_string())? | |
| 334 | + | .ok_or("pipeline missing")?; | |
| 335 | + | let pipeline = | |
| 336 | + | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; | |
| 337 | + | Ok(pipeline | |
| 338 | + | .artifacts | |
| 339 | + | .iter() | |
| 340 | + | .find(|a| a.name == name) | |
| 341 | + | .map(|a| ArtifactSpec { | |
| 342 | + | name: a.name.clone(), | |
| 343 | + | path: a.path.clone(), | |
| 344 | + | browse: a.browse, | |
| 345 | + | has_meta: !a.meta.is_empty(), | |
| 346 | + | })) | |
| 347 | + | } | |
| 348 | + | ||
| 349 | + | /// Store one uploaded artifact tar into the run's scratch directory. | |
| 350 | + | /// | |
| 351 | + | /// Scratch sits next to the artifacts' final home (same filesystem, so | |
| 352 | + | /// [`finish_run`]'s swap is a rename) and is only moved into place once the | |
| 353 | + | /// run reports in. | |
| 354 | + | pub fn store_upload( | |
| 355 | + | app: &App, | |
| 356 | + | run: &anvil_core::CiRun, | |
| 357 | + | spec: &ArtifactSpec, | |
| 358 | + | tar: &[u8], | |
| 359 | + | ) -> Result<Stored, String> { | |
| 360 | + | let scratch = scratch_dir(app, run); | |
| 361 | + | std::fs::create_dir_all(&scratch).map_err(|e| format!("creating scratch dir: {e}"))?; | |
| 362 | + | let (size, is_dir) = store_artifact(spec, tar, &scratch)?; | |
| 363 | + | Ok(Stored { size, is_dir }) | |
| 364 | + | } | |
| 365 | + | ||
| 366 | + | /// Where a run's artifacts accumulate before the swap. | |
| 367 | + | fn scratch_dir(app: &App, run: &anvil_core::CiRun) -> std::path::PathBuf { | |
| 368 | + | app.config | |
| 369 | + | .artifacts_dir() | |
| 370 | + | .join(run.repo_id.to_string()) | |
| 371 | + | .join(format!(".collecting-{}", run.id)) | |
| 372 | + | } | |
| 373 | + | ||
| 374 | + | /// Record a finished job: swap its artifacts into place, mask and store the | |
| 375 | + | /// log, set the status, and fire the deploy webhook if it qualifies. | |
| 376 | + | pub async fn finish_run(app: &App, run_id: i64, result: JobResult) -> Result<(), String> { | |
| 377 | + | let (owner, repo, repo_path, run) = resolve(app, run_id).await?; | |
| 378 | + | let scratch = scratch_dir(app, &run); | |
| 379 | + | let mut log = result.log; | |
| 380 | + | ||
| 381 | + | let (status, collected) = match result.runner_error { | |
| 382 | + | Some(e) => { | |
| 271 | 383 | log.push_str(&format!("\n[runner error] {e}\n")); | |
| 272 | 384 | (ci::status::ERROR, Vec::new()) | |
| 273 | 385 | } | |
| 386 | + | None if result.exit_code == 0 => (ci::status::SUCCESS, result.artifacts), | |
| 387 | + | None => { | |
| 388 | + | log.push_str(&format!("\n[exited with status {}]\n", result.exit_code)); | |
| 389 | + | (ci::status::FAILURE, result.artifacts) | |
| 390 | + | } | |
| 274 | 391 | }; | |
| 275 | 392 | ||
| 276 | 393 | // Swap the collected set into place, replacing any earlier run's | |
| 277 | 394 | // artifacts for this commit, then record the rows. | |
| 278 | 395 | if !collected.is_empty() { | |
| 396 | + | let artifacts_root = app.config.artifacts_dir(); | |
| 279 | 397 | let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit); | |
| 280 | 398 | let swap = async { | |
| 281 | 399 | ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit) | |
| ⋯ 32 unchanged lines | |||
| 314 | 432 | ||
| 315 | 433 | // A step that echoes its environment (`set -x`, `curl -v`, a failing | |
| 316 | 434 | // command that prints its arguments) would otherwise publish the value on | |
| 317 | - | // a page anyone with read access can see. | |
| 318 | - | mask_secrets(&mut log, &env); | |
| 435 | + | // a page anyone with read access can see. The values come from the lease | |
| 436 | + | // rather than the vault, which may have re-sealed while the job ran. | |
| 437 | + | mask_secrets(&mut log, &app.jobs.secrets_for(run_id)); | |
| 319 | 438 | ci::append_log(&app.db, run_id, &log).await.ok(); | |
| 320 | 439 | ci::finish(&app.db, run_id, status).await.ok(); | |
| 440 | + | app.jobs.release(run_id); | |
| 321 | 441 | tracing::info!("ci: run {run_id} {status}"); | |
| 322 | 442 | ||
| 323 | 443 | // Continuous deployment: on a green run of the configured deploy repo's | |
| 324 | 444 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — | |
| 325 | 445 | // no other repository can trigger it, even with passing CI. | |
| 326 | - | if status == ci::status::SUCCESS | |
| 327 | - | && app | |
| 328 | - | .config | |
| 329 | - | .ci | |
| 330 | - | .is_deploy_target(&owner.username, &repo.name, &run.ref_name) | |
| 446 | + | if status == ci::status::SUCCESS && app.config.ci.is_deploy_target(&owner, &repo, &run.ref_name) | |
| 331 | 447 | { | |
| 332 | - | deploy(app, &owner.username, &repo.name, &run).await; | |
| 448 | + | deploy(app, &owner, &repo, &run).await; | |
| 333 | 449 | } | |
| 334 | 450 | Ok(()) | |
| 335 | 451 | } | |
| 336 | 452 | ||
| 453 | + | /// Finish a run that never got far enough to produce a result. | |
| 454 | + | async fn fail_run(app: &App, run_id: i64, note: &str) { | |
| 455 | + | ci::append_log(&app.db, run_id, note).await.ok(); | |
| 456 | + | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); | |
| 457 | + | app.jobs.release(run_id); | |
| 458 | + | } | |
| 459 | + | ||
| 460 | + | /// (owner username, repo name, repo path, run) for a run id. | |
| 461 | + | async fn resolve( | |
| 462 | + | app: &App, | |
| 463 | + | run_id: i64, | |
| 464 | + | ) -> Result<(String, String, std::path::PathBuf, anvil_core::CiRun), String> { | |
| 465 | + | let run = ci::get(&app.db, run_id) | |
| 466 | + | .await | |
| 467 | + | .map_err(|e| e.to_string())? | |
| 468 | + | .ok_or("run not found")?; | |
| 469 | + | let repo = repos::find_by_id(&app.db, run.repo_id) | |
| 470 | + | .await | |
| 471 | + | .map_err(|e| e.to_string())? | |
| 472 | + | .ok_or("repository not found")?; | |
| 473 | + | let owner = users::find_by_id(&app.db, repo.owner_id) | |
| 474 | + | .await | |
| 475 | + | .map_err(|e| e.to_string())? | |
| 476 | + | .ok_or("owner not found")?; | |
| 477 | + | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); | |
| 478 | + | Ok((owner.username, repo.name, repo_path, run)) | |
| 479 | + | } | |
| 480 | + | ||
| 337 | 481 | /// Enforce `[ci] artifact_quota_mb` for one repository: while over budget, | |
| 338 | 482 | /// delete the oldest commit's artifacts (rows + directory). Branch-tip | |
| 339 | 483 | /// commits and the just-stored commit are pinned. Deterministic — runs after | |
| ⋯ 97 unchanged lines | |||
| 437 | 581 | } | |
| 438 | 582 | } | |
| 439 | 583 | ||
| 440 | - | /// Where a collected artifact's bytes go. | |
| 441 | - | /// | |
| 442 | - | /// Exists so the download loop can hand each tar off and drop it rather than | |
| 443 | - | /// buffering every artifact — `ci.artifact_run_max_mb` defaults to 512, which | |
| 444 | - | /// is not an amount to hold in RAM on the host anvil runs on. The in-process | |
| 445 | - | /// implementation writes to the scratch directory; a remote runner's uploads | |
| 446 | - | /// it. Returning [`Stored`] rather than `()` is what lets the caller charge | |
| 447 | - | /// the run budget the size that actually landed. | |
| 448 | - | #[async_trait::async_trait] | |
| 449 | - | pub trait ArtifactSink: Send { | |
| 450 | - | async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String>; | |
| 451 | - | } | |
| 452 | - | ||
| 453 | - | /// The in-process sink: straight into the run's scratch directory, which | |
| 454 | - | /// `process` then renames into the commit's artifact directory. | |
| 455 | - | struct ScratchSink<'a> { | |
| 456 | - | scratch: &'a Path, | |
| 457 | - | } | |
| 458 | - | ||
| 459 | - | #[async_trait::async_trait] | |
| 460 | - | impl ArtifactSink for ScratchSink<'_> { | |
| 461 | - | async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String> { | |
| 462 | - | // Created lazily rather than up front: a run whose artifacts all fail | |
| 463 | - | // to download should leave no empty scratch directory behind. | |
| 464 | - | std::fs::create_dir_all(self.scratch).map_err(|e| format!("creating scratch dir: {e}"))?; | |
| 465 | - | let (size, is_dir) = store_artifact(spec, tar, self.scratch)?; | |
| 466 | - | Ok(Stored { size, is_dir }) | |
| 467 | - | } | |
| 468 | - | } | |
| 469 | - | ||
| 470 | 584 | /// Replace every secret value in `log` with `***`. | |
| 471 | 585 | /// | |
| 472 | 586 | /// Only values worth hiding: very short ones (a one-character secret) would | |
| ⋯ 6 unchanged lines | |||
| 479 | 593 | } | |
| 480 | 594 | } | |
| 481 | 595 | ||
| 482 | - | /// Execute a job in a sandboxed container, streaming output into `log`. | |
| 483 | - | /// Returns the exit code and whatever artifacts `sink` accepted (none on | |
| 484 | - | /// timeout — the container is already gone). | |
| 485 | - | /// | |
| 486 | - | /// The job container never sees the Docker socket and gets no mounts of any | |
| 487 | - | /// kind (the checkout is *uploaded*, not bind-mounted; artifacts are | |
| 488 | - | /// *downloaded* out the same way). All capabilities are dropped and | |
| 489 | - | /// `no-new-privileges` is set unconditionally; pids/memory/cpu caps, the | |
| 490 | - | /// wall-clock timeout, network access and the container user come from | |
| 491 | - | /// `spec.sandbox`. The image allowlist was applied when the spec was built — | |
| 492 | - | /// a runner receiving a spec does not get to widen it. | |
| 493 | - | pub async fn execute( | |
| 494 | - | spec: &JobSpec, | |
| 495 | - | tar: Vec<u8>, | |
| 496 | - | log: &mut String, | |
| 497 | - | sink: &mut dyn ArtifactSink, | |
| 498 | - | ) -> Result<(i64, Vec<CollectedArtifact>), String> { | |
| 499 | - | let platform = spec.platform.clone().unwrap_or_default(); | |
| 500 | - | let docker = anvil_docker::connect()?; | |
| 501 | - | anvil_docker::ensure_image(&docker, &spec.image, &platform).await?; | |
| 502 | - | let sb = &spec.sandbox; | |
| 503 | - | ||
| 504 | - | // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap. | |
| 505 | - | let host_config = HostConfig { | |
| 506 | - | cap_drop: Some(vec!["ALL".to_string()]), | |
| 507 | - | security_opt: Some(vec!["no-new-privileges:true".to_string()]), | |
| 508 | - | pids_limit: (sb.pids_limit > 0).then_some(sb.pids_limit), | |
| 509 | - | memory: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024), | |
| 510 | - | memory_swap: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024), | |
| 511 | - | nano_cpus: (sb.cpus > 0.0).then_some((sb.cpus * 1e9) as i64), | |
| 512 | - | network_mode: (!sb.network).then(|| "none".to_string()), | |
| 513 | - | ..Default::default() | |
| 514 | - | }; | |
| 515 | - | let config = Config { | |
| 516 | - | image: Some(spec.image.clone()), | |
| 517 | - | cmd: Some(vec![ | |
| 518 | - | "sh".to_string(), | |
| 519 | - | "-c".to_string(), | |
| 520 | - | spec.script.clone(), | |
| 521 | - | ]), | |
| 522 | - | env: (!spec.env.is_empty()) | |
| 523 | - | .then(|| spec.env.iter().map(|(k, v)| format!("{k}={v}")).collect()), | |
| 524 | - | working_dir: Some(WORKDIR.to_string()), | |
| 525 | - | user: (!sb.run_as.is_empty()).then(|| sb.run_as.clone()), | |
| 526 | - | host_config: Some(host_config), | |
| 527 | - | ..Default::default() | |
| 528 | - | }; | |
| 529 | - | // An explicit platform needs the options struct; without one, pass None so | |
| 530 | - | // the daemon picks its native architecture exactly as before. | |
| 531 | - | let opts = spec.platform.as_ref().map(|p| CreateContainerOptions { | |
| 532 | - | name: String::new(), | |
| 533 | - | platform: Some(p.clone()), | |
| 534 | - | }); | |
| 535 | - | let created = docker | |
| 536 | - | .create_container(opts, config) | |
| 537 | - | .await | |
| 538 | - | .map_err(|e| format!("create container: {e}"))?; | |
| 539 | - | let id = created.id; | |
| 540 | - | ||
| 541 | - | // Upload the checkout (tar entries are under `workspace/`, extracted at `/`). | |
| 542 | - | docker | |
| 543 | - | .upload_to_container( | |
| 544 | - | &id, | |
| 545 | - | Some(UploadToContainerOptions { | |
| 546 | - | path: "/".to_string(), | |
| 547 | - | ..Default::default() | |
| 548 | - | }), | |
| 549 | - | tar.into(), | |
| 550 | - | ) | |
| 551 | - | .await | |
| 552 | - | .map_err(|e| format!("upload checkout: {e}"))?; | |
| 553 | - | ||
| 554 | - | docker | |
| 555 | - | .start_container(&id, None::<StartContainerOptions<String>>) | |
| 556 | - | .await | |
| 557 | - | .map_err(|e| format!("start container: {e}"))?; | |
| 558 | - | ||
| 559 | - | // Stream logs and wait for the exit code, bounded by the wall-clock | |
| 560 | - | // timeout. The container is force-removed on every path (which also kills | |
| 561 | - | // a still-running job after a timeout). | |
| 562 | - | let run = async { | |
| 563 | - | let mut logs = docker.logs( | |
| 564 | - | &id, | |
| 565 | - | Some(LogsOptions::<String> { | |
| 566 | - | follow: true, | |
| 567 | - | stdout: true, | |
| 568 | - | stderr: true, | |
| 569 | - | ..Default::default() | |
| 570 | - | }), | |
| 571 | - | ); | |
| 572 | - | while let Some(item) = logs.next().await { | |
| 573 | - | match item { | |
| 574 | - | Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())), | |
| 575 | - | Err(e) => { | |
| 576 | - | log.push_str(&format!("\n[log stream error] {e}\n")); | |
| 577 | - | break; | |
| 578 | - | } | |
| 579 | - | } | |
| 580 | - | } | |
| 581 | - | ||
| 582 | - | // Non-zero exit codes surface as a wait error in bollard. | |
| 583 | - | let mut code = 0i64; | |
| 584 | - | let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>); | |
| 585 | - | while let Some(item) = wait.next().await { | |
| 586 | - | match item { | |
| 587 | - | Ok(resp) => code = resp.status_code, | |
| 588 | - | Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c, | |
| 589 | - | Err(e) => return Err(format!("wait: {e}")), | |
| 590 | - | } | |
| 591 | - | } | |
| 592 | - | Ok(code) | |
| 593 | - | }; | |
| 594 | - | let result = match sb.timeout_secs { | |
| 595 | - | 0 => run.await, | |
| 596 | - | secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run) | |
| 597 | - | .await | |
| 598 | - | .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))), | |
| 599 | - | }; | |
| 600 | - | ||
| 601 | - | // Artifacts come out of the (now stopped) container before it is removed. | |
| 602 | - | let collected = match &result { | |
| 603 | - | Ok(_) if !spec.artifacts.is_empty() => { | |
| 604 | - | collect_artifacts(&docker, &id, spec, log, sink).await | |
| 605 | - | } | |
| 606 | - | _ => Vec::new(), | |
| 607 | - | }; | |
| 608 | - | ||
| 609 | - | let _ = docker | |
| 610 | - | .remove_container( | |
| 611 | - | &id, | |
| 612 | - | Some(RemoveContainerOptions { | |
| 613 | - | force: true, | |
| 614 | - | ..Default::default() | |
| 615 | - | }), | |
| 616 | - | ) | |
| 617 | - | .await; | |
| 618 | - | ||
| 619 | - | result.map(|code| (code, collected)) | |
| 620 | - | } | |
| 621 | - | ||
| 622 | - | /// Collect the job's declared artifacts from the stopped container, handing | |
| 623 | - | /// each tar to `sink` as it is downloaded. Failures are per-artifact: each is | |
| 624 | - | /// logged and skipped, never failing the run. | |
| 625 | - | /// | |
| 626 | - | /// One artifact is in memory at a time by construction — the tar is dropped | |
| 627 | - | /// once the sink has taken it — which is what keeps `artifact_run_max_mb` | |
| 628 | - | /// (512 MiB by default) a disk budget rather than a memory one. | |
| 629 | - | async fn collect_artifacts( | |
| 630 | - | docker: &Docker, | |
| 631 | - | id: &str, | |
| 632 | - | job: &JobSpec, | |
| 633 | - | log: &mut String, | |
| 634 | - | sink: &mut dyn ArtifactSink, | |
| 635 | - | ) -> Vec<CollectedArtifact> { | |
| 636 | - | // Extractor outputs first: artifact name → key → value. | |
| 637 | - | let mut metas: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); | |
| 638 | - | if job.artifacts.iter().any(|a| a.has_meta) { | |
| 639 | - | match download_tar(docker, id, META_DIR, META_TAR_CAP).await { | |
| 640 | - | Ok(Some(bytes)) => metas = parse_meta_tar(&bytes), | |
| 641 | - | Ok(None) => log.push_str("\n[artifacts: extractor output exceeded its cap]\n"), | |
| 642 | - | Err(e) => log.push_str(&format!("\n[artifacts: reading extractor output: {e}]\n")), | |
| 643 | - | } | |
| 644 | - | } | |
| 645 | - | ||
| 646 | - | let per_artifact_cap = mb_cap(job.sandbox.artifact_max_mb); | |
| 647 | - | let mut run_budget = mb_cap(job.sandbox.artifact_run_max_mb); | |
| 648 | - | let mut collected = Vec::new(); | |
| 649 | - | for spec in &job.artifacts { | |
| 650 | - | let cap = per_artifact_cap.min(run_budget); | |
| 651 | - | let note = |log: &mut String, what: &str| { | |
| 652 | - | log.push_str(&format!("\n[artifact {}: {what}]\n", spec.name)); | |
| 653 | - | }; | |
| 654 | - | let bytes = match download_tar(docker, id, &format!("{WORKDIR}/{}", spec.path), cap).await { | |
| 655 | - | Ok(Some(bytes)) => bytes, | |
| 656 | - | Ok(None) => { | |
| 657 | - | note(log, "exceeds the size cap; skipped"); | |
| 658 | - | continue; | |
| 659 | - | } | |
| 660 | - | Err(e) => { | |
| 661 | - | note(log, &format!("download failed ({e}); skipped")); | |
| 662 | - | continue; | |
| 663 | - | } | |
| 664 | - | }; | |
| 665 | - | match sink.put(spec, &bytes).await { | |
| 666 | - | Ok(Stored { size, is_dir }) => { | |
| 667 | - | run_budget = run_budget.saturating_sub(size as u64); | |
| 668 | - | let meta = metas.get(&spec.name).cloned().unwrap_or_default(); | |
| 669 | - | collected.push(CollectedArtifact { | |
| 670 | - | name: spec.name.clone(), | |
| 671 | - | size, | |
| 672 | - | is_dir, | |
| 673 | - | browse: spec.browse, | |
| 674 | - | meta: serde_json::to_string(&meta).unwrap_or_else(|_| "{}".into()), | |
| 675 | - | }); | |
| 676 | - | } | |
| 677 | - | Err(e) => note(log, &format!("storing failed ({e}); skipped")), | |
| 678 | - | } | |
| 679 | - | } | |
| 680 | - | collected | |
| 681 | - | } | |
| 682 | - | ||
| 683 | - | /// `0` (unlimited) → `u64::MAX`, otherwise MiB → bytes. | |
| 684 | - | fn mb_cap(mb: i64) -> u64 { | |
| 685 | - | if mb <= 0 { | |
| 686 | - | u64::MAX | |
| 687 | - | } else { | |
| 688 | - | mb as u64 * 1024 * 1024 | |
| 689 | - | } | |
| 690 | - | } | |
| 691 | - | ||
| 692 | - | /// Download `path` from the container as a tar, buffering at most `cap` bytes | |
| 693 | - | /// (`Ok(None)` when exceeded). | |
| 694 | - | async fn download_tar( | |
| 695 | - | docker: &Docker, | |
| 696 | - | id: &str, | |
| 697 | - | path: &str, | |
| 698 | - | cap: u64, | |
| 699 | - | ) -> Result<Option<Vec<u8>>, String> { | |
| 700 | - | let mut stream = docker.download_from_container( | |
| 701 | - | id, | |
| 702 | - | Some(DownloadFromContainerOptions { | |
| 703 | - | path: path.to_string(), | |
| 704 | - | }), | |
| 705 | - | ); | |
| 706 | - | let mut buf = Vec::new(); | |
| 707 | - | while let Some(chunk) = stream.next().await { | |
| 708 | - | let chunk = chunk.map_err(|e| e.to_string())?; | |
| 709 | - | if (buf.len() + chunk.len()) as u64 > cap { | |
| 710 | - | return Ok(None); | |
| 711 | - | } | |
| 712 | - | buf.extend_from_slice(&chunk); | |
| 713 | - | } | |
| 714 | - | Ok(Some(buf)) | |
| 715 | - | } | |
| 716 | - | ||
| 717 | - | /// Parse the extractor-output tar (`anvil-meta/<artifact>/<key>` files) into | |
| 718 | - | /// artifact → key → trimmed value. | |
| 719 | - | fn parse_meta_tar(bytes: &[u8]) -> BTreeMap<String, BTreeMap<String, String>> { | |
| 720 | - | let mut out: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); | |
| 721 | - | let mut archive = tar::Archive::new(bytes); | |
| 722 | - | let Ok(entries) = archive.entries() else { | |
| 723 | - | return out; | |
| 724 | - | }; | |
| 725 | - | for entry in entries.flatten() { | |
| 726 | - | if !entry.header().entry_type().is_file() { | |
| 727 | - | continue; | |
| 728 | - | } | |
| 729 | - | let Ok(path) = entry.path() else { continue }; | |
| 730 | - | // anvil-meta/<artifact>/<key> | |
| 731 | - | let parts: Vec<String> = path | |
| 732 | - | .components() | |
| 733 | - | .skip(1) | |
| 734 | - | .map(|c| c.as_os_str().to_string_lossy().into_owned()) | |
| 735 | - | .collect(); | |
| 736 | - | let [artifact, key] = parts.as_slice() else { | |
| 737 | - | continue; | |
| 738 | - | }; | |
| 739 | - | let (artifact, key) = (artifact.clone(), key.clone()); | |
| 740 | - | let mut value = String::new(); | |
| 741 | - | let _ = entry.take(META_VALUE_CAP as u64).read_to_string(&mut value); | |
| 742 | - | let value = value.trim().to_string(); | |
| 743 | - | if !value.is_empty() { | |
| 744 | - | out.entry(artifact).or_default().insert(key, value); | |
| 745 | - | } | |
| 746 | - | } | |
| 747 | - | out | |
| 748 | - | } | |
| 749 | - | ||
| 750 | 596 | /// Write one downloaded artifact tar into `scratch`, returning (size, is_dir). | |
| 751 | 597 | /// | |
| 752 | 598 | /// - a file is stored as-is at `scratch/<name>` | |
| ⋯ 178 unchanged lines | |||
| 931 | 777 | let root = dir.path().join("doc"); | |
| 932 | 778 | assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>"); | |
| 933 | 779 | assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>"); | |
| 934 | - | } | |
| 935 | - | ||
| 936 | - | #[test] | |
| 937 | - | fn parses_meta_tar_with_trimmed_capped_values() { | |
| 938 | - | let tar = tar_of(&[ | |
| 939 | - | ("anvil-meta", None), | |
| 940 | - | ("anvil-meta/bin", None), | |
| 941 | - | ("anvil-meta/bin/version", Some("anvild 0.0.0\n")), | |
| 942 | - | ("anvil-meta/bin/empty", Some(" \n")), | |
| 943 | - | ("anvil-meta/doc", None), | |
| 944 | - | ("anvil-meta/doc/pages", Some("42")), | |
| 945 | - | ]); | |
| 946 | - | let metas = parse_meta_tar(&tar); | |
| 947 | - | assert_eq!(metas["bin"]["version"], "anvild 0.0.0"); | |
| 948 | - | assert_eq!(metas["doc"]["pages"], "42"); | |
| 949 | - | assert!(!metas["bin"].contains_key("empty"), "blank values dropped"); | |
| 950 | 780 | } | |
| 951 | 781 | } | |
modifiedcrates/anvil-cli/src/main.rs+1 −1
| ⋯ 161 unchanged lines | |||
| 162 | 162 | // through `app.ci_tx` (set here so handlers can notify it). | |
| 163 | 163 | let (ci_tx, ci_rx) = tokio::sync::mpsc::unbounded_channel(); | |
| 164 | 164 | app.ci_tx = Some(ci_tx); | |
| 165 | - | tokio::spawn(anvil_ci::run_worker(app.clone(), ci_rx)); | |
| 165 | + | tokio::spawn(anvil_ci::run_dispatcher(app.clone(), ci_rx)); | |
| 166 | 166 | ||
| 167 | 167 | // Agent sessions outlive the process that started them: the container is | |
| 168 | 168 | // the daemon's, not ours. Reconcile before anything can create more, so a | |
| ⋯ 179 unchanged lines | |||
modifiedcrates/anvil-core/Cargo.toml+1 −0
| ⋯ 7 unchanged lines | |||
| 8 | 8 | description = "Domain model, persistence, and on-disk repository storage for the anvil git forge." | |
| 9 | 9 | ||
| 10 | 10 | [dependencies] | |
| 11 | + | anvil-job.workspace = true | |
| 11 | 12 | async-trait.workspace = true | |
| 12 | 13 | gix.workspace = true | |
| 13 | 14 | toasty.workspace = true | |
| ⋯ 27 unchanged lines | |||
modifiedcrates/anvil-core/src/ci.rs+15 −0
| ⋯ 255 unchanged lines | |||
| 256 | 256 | Ok(ids) | |
| 257 | 257 | } | |
| 258 | 258 | ||
| 259 | + | /// Put one run back on the queue. | |
| 260 | + | /// | |
| 261 | + | /// Used when a runner's lease expires: the job may still be executing on an | |
| 262 | + | /// unreachable machine, but from anvil's side it is no longer accounted for, | |
| 263 | + | /// and leaving it `running` forever is worse than the chance of a second | |
| 264 | + | /// attempt. | |
| 265 | + | pub async fn requeue(db: &toasty::Db, id: i64) -> Result<()> { | |
| 266 | + | let mut conn = db.clone(); | |
| 267 | + | let Some(mut run) = get(db, id).await? else { | |
| 268 | + | return Ok(()); | |
| 269 | + | }; | |
| 270 | + | run.update().status(status::QUEUED).exec(&mut conn).await?; | |
| 271 | + | Ok(()) | |
| 272 | + | } | |
| 273 | + | ||
| 259 | 274 | /// List all queued run ids (oldest first) — used on startup to drain the queue. | |
| 260 | 275 | pub async fn queued_ids(db: &toasty::Db) -> Result<Vec<i64>> { | |
| 261 | 276 | let mut conn = db.clone(); | |
| ⋯ 251 unchanged lines | |||
modifiedcrates/anvil-core/src/config.rs+12 −0
| ⋯ 135 unchanged lines | |||
| 136 | 136 | #[derive(Clone, Debug, Deserialize, Serialize)] | |
| 137 | 137 | #[serde(default)] | |
| 138 | 138 | pub struct CiConfig { | |
| 139 | + | /// Shared secret a runner presents as `X-Anvil-Runner-Token` to claim and | |
| 140 | + | /// report jobs (see `docs/remote-runners.md`). **Empty refuses every | |
| 141 | + | /// runner**, which means no CI runs at all — anvil does not execute jobs | |
| 142 | + | /// itself any more. | |
| 143 | + | /// | |
| 144 | + | /// A config-file secret rather than a credential type of its own, matching | |
| 145 | + | /// [`deploy_secret`](CiConfig::deploy_secret): API tokens are read-only and | |
| 146 | + | /// Bearer-only on GET/HEAD, and a runner must POST. Per-runner DB-backed | |
| 147 | + | /// tokens are the right end state; one shared secret is enough for a | |
| 148 | + | /// single-tenant forge and needs no token-management UI. | |
| 149 | + | pub runner_token: String, | |
| 139 | 150 | /// The one repository (`owner/name`) permitted to trigger the deploy | |
| 140 | 151 | /// webhook. Empty disables deploys entirely. | |
| 141 | 152 | pub deploy_repo: String, | |
| ⋯ 190 unchanged lines | |||
| 332 | 343 | impl Default for CiConfig { | |
| 333 | 344 | fn default() -> Self { | |
| 334 | 345 | Self { | |
| 346 | + | runner_token: String::new(), | |
| 335 | 347 | deploy_repo: String::new(), | |
| 336 | 348 | deploy_webhook: String::new(), | |
| 337 | 349 | deploy_secret: String::new(), | |
| ⋯ 347 unchanged lines | |||
addedcrates/anvil-core/src/jobs.rs+168 −0
| 1 | + | //! Dispatch state for remote job runners (see `docs/remote-runners.md`). | |
| 2 | + | //! | |
| 3 | + | //! anvil no longer executes CI itself; it hands jobs to runners that dial in | |
| 4 | + | //! and claim them. This holds the two pieces of state that implies: who is | |
| 5 | + | //! currently running what, and a way to wake a parked claim request when work | |
| 6 | + | //! arrives. | |
| 7 | + | //! | |
| 8 | + | //! **Deliberately in memory**, not columns on `CiRun`. Toasty has no migration | |
| 9 | + | //! support yet (see `TODO.md` and DEPLOY.md's Operations note), so adding | |
| 10 | + | //! `claimed_by`/`lease_expires_at` to the model would need a hand-written | |
| 11 | + | //! migration against the live database. It also costs nothing to keep here: | |
| 12 | + | //! anvild is the only dispatcher, so a lease has no reason to outlive it, and | |
| 13 | + | //! the case where it dies mid-job is already covered — `ci::requeue_interrupted` | |
| 14 | + | //! re-queues everything still marked `running` at startup. | |
| 15 | + | //! | |
| 16 | + | //! Same shape as [`crate::agent::Registry`], for the same reason. | |
| 17 | + | ||
| 18 | + | use std::{ | |
| 19 | + | collections::HashMap, | |
| 20 | + | sync::{ | |
| 21 | + | Arc, | |
| 22 | + | Mutex, | |
| 23 | + | }, | |
| 24 | + | time::{ | |
| 25 | + | Duration, | |
| 26 | + | Instant, | |
| 27 | + | }, | |
| 28 | + | }; | |
| 29 | + | ||
| 30 | + | pub use anvil_job::{ | |
| 31 | + | HEARTBEAT_INTERVAL, | |
| 32 | + | LEASE_TTL, | |
| 33 | + | }; | |
| 34 | + | ||
| 35 | + | /// Who holds a run, and until when. | |
| 36 | + | #[derive(Clone, Debug)] | |
| 37 | + | pub struct Lease { | |
| 38 | + | pub runner: String, | |
| 39 | + | pub expires_at: Instant, | |
| 40 | + | /// The secret values handed to this job, kept so the log can be masked | |
| 41 | + | /// when the result comes back. | |
| 42 | + | /// | |
| 43 | + | /// Stashed rather than re-read from the vault at finish time, because | |
| 44 | + | /// `Vault::take` fails once the repository's unlock TTL lapses — and a job | |
| 45 | + | /// may well outlive it. Re-reading would mean a run that took longer than | |
| 46 | + | /// the unlock silently skips masking, which is exactly the case where the | |
| 47 | + | /// log is most likely to contain something. No new exposure: these values | |
| 48 | + | /// are already in this process's vault. | |
| 49 | + | pub secrets: Vec<(String, String)>, | |
| 50 | + | } | |
| 51 | + | ||
| 52 | + | /// Shared dispatch state. Cheap to clone; all clones share one map. | |
| 53 | + | #[derive(Clone, Default)] | |
| 54 | + | pub struct Dispatch { | |
| 55 | + | inner: Arc<Inner>, | |
| 56 | + | } | |
| 57 | + | ||
| 58 | + | #[derive(Default)] | |
| 59 | + | struct Inner { | |
| 60 | + | leases: Mutex<HashMap<i64, Lease>>, | |
| 61 | + | wake: tokio::sync::Notify, | |
| 62 | + | } | |
| 63 | + | ||
| 64 | + | impl Dispatch { | |
| 65 | + | pub fn new() -> Self { | |
| 66 | + | Self::default() | |
| 67 | + | } | |
| 68 | + | ||
| 69 | + | /// Wake every parked claim request. Called when a run is enqueued. | |
| 70 | + | /// | |
| 71 | + | /// `notify_waiters` rather than `notify_one`: a runner that wakes and finds | |
| 72 | + | /// the queue already emptied by another simply parks again, whereas | |
| 73 | + | /// `notify_one` can hand the permit to a runner that is about to give up. | |
| 74 | + | pub fn wake(&self) { | |
| 75 | + | self.inner.wake.notify_waiters(); | |
| 76 | + | } | |
| 77 | + | ||
| 78 | + | /// Park until there might be work, or until `timeout` elapses. | |
| 79 | + | pub async fn wait_for_work(&self, timeout: Duration) { | |
| 80 | + | let _ = tokio::time::timeout(timeout, self.inner.wake.notified()).await; | |
| 81 | + | } | |
| 82 | + | ||
| 83 | + | /// Record that `runner` holds `run_id`, stashing the job's secrets for | |
| 84 | + | /// masking when the result arrives. | |
| 85 | + | pub fn claim(&self, run_id: i64, runner: &str, secrets: Vec<(String, String)>) { | |
| 86 | + | self.inner.leases.lock().unwrap().insert( | |
| 87 | + | run_id, | |
| 88 | + | Lease { | |
| 89 | + | runner: runner.to_string(), | |
| 90 | + | expires_at: Instant::now() + LEASE_TTL, | |
| 91 | + | secrets, | |
| 92 | + | }, | |
| 93 | + | ); | |
| 94 | + | } | |
| 95 | + | ||
| 96 | + | /// The secrets handed to `run_id`, for masking its log. | |
| 97 | + | pub fn secrets_for(&self, run_id: i64) -> Vec<(String, String)> { | |
| 98 | + | self.inner | |
| 99 | + | .leases | |
| 100 | + | .lock() | |
| 101 | + | .unwrap() | |
| 102 | + | .get(&run_id) | |
| 103 | + | .map(|l| l.secrets.clone()) | |
| 104 | + | .unwrap_or_default() | |
| 105 | + | } | |
| 106 | + | ||
| 107 | + | /// Extend `runner`'s lease on `run_id`. False when it holds no such lease, | |
| 108 | + | /// which is how a runner learns its job was requeued out from under it. | |
| 109 | + | pub fn heartbeat(&self, run_id: i64, runner: &str) -> bool { | |
| 110 | + | let mut leases = self.inner.leases.lock().unwrap(); | |
| 111 | + | match leases.get_mut(&run_id) { | |
| 112 | + | Some(lease) if lease.runner == runner => { | |
| 113 | + | lease.expires_at = Instant::now() + LEASE_TTL; | |
| 114 | + | true | |
| 115 | + | } | |
| 116 | + | _ => false, | |
| 117 | + | } | |
| 118 | + | } | |
| 119 | + | ||
| 120 | + | /// Whether `runner` currently holds `run_id`. The authorization check for | |
| 121 | + | /// every per-job endpoint: a valid token gets you a job, but only the | |
| 122 | + | /// holder may fetch its checkout or post its result. | |
| 123 | + | pub fn holds(&self, run_id: i64, runner: &str) -> bool { | |
| 124 | + | self.inner | |
| 125 | + | .leases | |
| 126 | + | .lock() | |
| 127 | + | .unwrap() | |
| 128 | + | .get(&run_id) | |
| 129 | + | .is_some_and(|l| l.runner == runner) | |
| 130 | + | } | |
| 131 | + | ||
| 132 | + | /// Whether anyone holds `run_id`. Lets the claim loop skip a run another | |
| 133 | + | /// runner took between the queue read and here. | |
| 134 | + | pub fn holds_any(&self, run_id: i64) -> bool { | |
| 135 | + | self.inner.leases.lock().unwrap().contains_key(&run_id) | |
| 136 | + | } | |
| 137 | + | ||
| 138 | + | /// Drop the lease on `run_id` (the job finished, one way or another). | |
| 139 | + | pub fn release(&self, run_id: i64) { | |
| 140 | + | self.inner.leases.lock().unwrap().remove(&run_id); | |
| 141 | + | } | |
| 142 | + | ||
| 143 | + | /// Remove and return every lease past its TTL. The caller requeues them. | |
| 144 | + | pub fn take_expired(&self) -> Vec<(i64, String)> { | |
| 145 | + | let now = Instant::now(); | |
| 146 | + | let mut leases = self.inner.leases.lock().unwrap(); | |
| 147 | + | let dead: Vec<_> = leases | |
| 148 | + | .iter() | |
| 149 | + | .filter(|(_, l)| l.expires_at <= now) | |
| 150 | + | .map(|(id, l)| (*id, l.runner.clone())) | |
| 151 | + | .collect(); | |
| 152 | + | for (id, _) in &dead { | |
| 153 | + | leases.remove(id); | |
| 154 | + | } | |
| 155 | + | dead | |
| 156 | + | } | |
| 157 | + | ||
| 158 | + | /// Run ids currently claimed, for the admin view and for logging. | |
| 159 | + | pub fn in_flight(&self) -> Vec<(i64, String)> { | |
| 160 | + | self.inner | |
| 161 | + | .leases | |
| 162 | + | .lock() | |
| 163 | + | .unwrap() | |
| 164 | + | .iter() | |
| 165 | + | .map(|(id, l)| (*id, l.runner.clone())) | |
| 166 | + | .collect() | |
| 167 | + | } | |
| 168 | + | } |
modifiedcrates/anvil-core/src/lib.rs+7 −0
| ⋯ 14 unchanged lines | |||
| 15 | 15 | pub mod db; | |
| 16 | 16 | pub mod error; | |
| 17 | 17 | pub mod issues; | |
| 18 | + | pub mod jobs; | |
| 18 | 19 | pub mod language; | |
| 19 | 20 | pub mod models; | |
| 20 | 21 | pub mod periodic; | |
| ⋯ 58 unchanged lines | |||
| 79 | 80 | /// restart, which is why startup reconciles rows against containers — see | |
| 80 | 81 | /// [`agent::Registry`]. | |
| 81 | 82 | pub sessions: agent::Registry, | |
| 83 | + | /// Which runner holds which CI run, plus the wakeup parked claim requests | |
| 84 | + | /// block on. In memory and lost on restart, which is safe because | |
| 85 | + | /// [`ci::requeue_interrupted`] re-queues anything still `running` at | |
| 86 | + | /// startup — see [`jobs::Dispatch`]. | |
| 87 | + | pub jobs: jobs::Dispatch, | |
| 82 | 88 | /// Server-wide secret keying CSRF tokens. Persisted in the data dir so | |
| 83 | 89 | /// tokens survive restarts. Wrapped in `Arc` to keep `App: Clone` cheap. | |
| 84 | 90 | csrf_secret: std::sync::Arc<[u8; 32]>, | |
| ⋯ 16 unchanged lines | |||
| 101 | 107 | vault: secrets::Vault::default(), | |
| 102 | 108 | user_vault: secrets::Vault::default(), | |
| 103 | 109 | sessions: agent::Registry::default(), | |
| 110 | + | jobs: jobs::Dispatch::new(), | |
| 104 | 111 | csrf_secret, | |
| 105 | 112 | }) | |
| 106 | 113 | } | |
| ⋯ 52 unchanged lines | |||
modifiedcrates/anvil-job/src/lib.rs+27 −0
| ⋯ 10 unchanged lines | |||
| 11 | 11 | //! never parses a pipeline. That keeps this format stable as the pipeline | |
| 12 | 12 | //! schema grows. | |
| 13 | 13 | ||
| 14 | + | use std::time::Duration; | |
| 15 | + | ||
| 14 | 16 | use serde::{ | |
| 15 | 17 | Deserialize, | |
| 16 | 18 | Serialize, | |
| ⋯ 16 unchanged lines | |||
| 33 | 35 | /// Per-value cap on extractor output, in bytes (after trimming). | |
| 34 | 36 | pub const META_VALUE_CAP: usize = 1024; | |
| 35 | 37 | ||
| 38 | + | /// How long a claim survives without a heartbeat before anvil requeues the run. | |
| 39 | + | /// | |
| 40 | + | /// Generous relative to [`HEARTBEAT_INTERVAL`] so a slow network or a busy | |
| 41 | + | /// runner does not lose a job that is running fine. It bounds how long a dead | |
| 42 | + | /// runner strands a run, not how long a job may take — a live runner holds its | |
| 43 | + | /// lease across a 30-minute build by heartbeating through it. | |
| 44 | + | pub const LEASE_TTL: Duration = Duration::from_secs(120); | |
| 45 | + | ||
| 46 | + | /// How often a runner heartbeats while a job runs. Here rather than in the | |
| 47 | + | /// runner so both halves of the pair stay in view of each other. | |
| 48 | + | pub const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(30); | |
| 49 | + | ||
| 50 | + | /// A MiB cap from config as a byte count: `0` (unlimited) → [`u64::MAX`]. | |
| 51 | + | /// | |
| 52 | + | /// Here rather than on either side because both apply the same caps — the | |
| 53 | + | /// runner while downloading artifacts, anvil while enforcing the per-repo | |
| 54 | + | /// quota — and they must agree on what zero means. | |
| 55 | + | pub fn mb_cap(mb: i64) -> u64 { | |
| 56 | + | if mb <= 0 { | |
| 57 | + | u64::MAX | |
| 58 | + | } else { | |
| 59 | + | (mb as u64).saturating_mul(1024 * 1024) | |
| 60 | + | } | |
| 61 | + | } | |
| 62 | + | ||
| 36 | 63 | /// What a runner says about itself when it asks for work. | |
| 37 | 64 | #[derive(Clone, Debug, Serialize, Deserialize)] | |
| 38 | 65 | pub struct RunnerInfo { | |
| ⋯ 106 unchanged lines | |||
modifiedcrates/anvil-web/Cargo.toml+2 −0
| ⋯ 8 unchanged lines | |||
| 9 | 9 | ||
| 10 | 10 | [dependencies] | |
| 11 | 11 | anvil-agent.workspace = true | |
| 12 | + | anvil-ci.workspace = true | |
| 12 | 13 | anvil-core.workspace = true | |
| 14 | + | anvil-job.workspace = true | |
| 13 | 15 | anvil-git.workspace = true | |
| 14 | 16 | axum.workspace = true | |
| 15 | 17 | axum-extra.workspace = true | |
| ⋯ 31 unchanged lines | |||
modifiedcrates/anvil-web/src/lib.rs+2 −0
| ⋯ 27 unchanged lines | |||
| 28 | 28 | pub mod git_http; | |
| 29 | 29 | pub mod oidc; | |
| 30 | 30 | pub mod pages; | |
| 31 | + | pub mod runner; | |
| 31 | 32 | pub mod secrets; | |
| 32 | 33 | pub mod todomd; | |
| 33 | 34 | pub mod ui; | |
| ⋯ 15 unchanged lines | |||
| 49 | 50 | ); // uploaded image attachments | |
| 50 | 51 | router = oidc::routes(router); // single sign-on, when configured | |
| 51 | 52 | router = secrets::routes(router); // sealed per-repo secrets + unlock API | |
| 53 | + | router = runner::routes(router); // job runners claiming and reporting CI | |
| 52 | 54 | router = git_http::routes(router); // smart-HTTP git endpoints | |
| 53 | 55 | router | |
| 54 | 56 | // Derives the per-request CSRF token so the layout can attach it to | |
| ⋯ 22 unchanged lines | |||
addedcrates/anvil-web/src/runner.rs+227 −0
| 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 | + | ||
| 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. | |
| 134 | + | const CLAIM_POLL: std::time::Duration = std::time::Duration::from_secs(55); | |
| 135 | + | ||
| 136 | + | /// The checkout for a claimed job, as an uncompressed tar. | |
| 137 | + | async 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 | + | ||
| 152 | + | async 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. | |
| 168 | + | async 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. | |
| 202 | + | async 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. | |
| 217 | + | fn 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 | + | } |
addedcrates/anvil-worker/Cargo.toml+28 −0
| 1 | + | [package] | |
| 2 | + | name = "anvil-worker" | |
| 3 | + | version.workspace = true | |
| 4 | + | edition.workspace = true | |
| 5 | + | license.workspace = true | |
| 6 | + | repository.workspace = true | |
| 7 | + | rust-version.workspace = true | |
| 8 | + | description = "anvil's job runner: dials out to a forge, claims CI jobs, and runs them in sandboxed containers on its own Docker daemon." | |
| 9 | + | ||
| 10 | + | # Named `anvil-worker` rather than `anvil-runner` because `anvil-runner:latest` | |
| 11 | + | # is already the *image* CI jobs and agent sessions run inside. | |
| 12 | + | [[bin]] | |
| 13 | + | name = "anvil-worker" | |
| 14 | + | path = "src/main.rs" | |
| 15 | + | ||
| 16 | + | [dependencies] | |
| 17 | + | anvil-docker.workspace = true | |
| 18 | + | anvil-job.workspace = true | |
| 19 | + | async-trait.workspace = true | |
| 20 | + | bollard.workspace = true | |
| 21 | + | clap = { workspace = true, features = ["env"] } | |
| 22 | + | tar.workspace = true | |
| 23 | + | futures-util.workspace = true | |
| 24 | + | reqwest.workspace = true | |
| 25 | + | serde_json.workspace = true | |
| 26 | + | tokio.workspace = true | |
| 27 | + | tracing.workspace = true | |
| 28 | + | tracing-subscriber.workspace = true |
addedcrates/anvil-worker/src/client.rs+187 −0
| 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 | + | } |
addedcrates/anvil-worker/src/executor.rs+364 −0
| 1 | + | //! Running one job in a sandboxed container. | |
| 2 | + | //! | |
| 3 | + | //! Lifted verbatim out of `anvil-ci`, where it ran in anvild's own process | |
| 4 | + | //! against the host's Docker socket. The sandbox is unchanged — that is the | |
| 5 | + | //! point: moving execution to another machine must not quietly relax it. | |
| 6 | + | //! | |
| 7 | + | //! The job container never sees the Docker socket and gets no mounts of any | |
| 8 | + | //! kind. The checkout is *uploaded* through the Docker API and artifacts are | |
| 9 | + | //! *downloaded* back out the same way, so a job cannot reach the runner's | |
| 10 | + | //! filesystem any more than it could reach anvil's. | |
| 11 | + | ||
| 12 | + | use std::{ | |
| 13 | + | collections::BTreeMap, | |
| 14 | + | io::Read, | |
| 15 | + | }; | |
| 16 | + | ||
| 17 | + | use anvil_job::{ | |
| 18 | + | ArtifactSpec, | |
| 19 | + | CollectedArtifact, | |
| 20 | + | JobSpec, | |
| 21 | + | META_DIR, | |
| 22 | + | META_TAR_CAP, | |
| 23 | + | META_VALUE_CAP, | |
| 24 | + | Stored, | |
| 25 | + | WORKDIR, | |
| 26 | + | mb_cap, | |
| 27 | + | }; | |
| 28 | + | use bollard::{ | |
| 29 | + | Docker, | |
| 30 | + | container::{ | |
| 31 | + | Config, | |
| 32 | + | CreateContainerOptions, | |
| 33 | + | DownloadFromContainerOptions, | |
| 34 | + | LogsOptions, | |
| 35 | + | RemoveContainerOptions, | |
| 36 | + | StartContainerOptions, | |
| 37 | + | UploadToContainerOptions, | |
| 38 | + | WaitContainerOptions, | |
| 39 | + | }, | |
| 40 | + | models::HostConfig, | |
| 41 | + | }; | |
| 42 | + | use futures_util::StreamExt; | |
| 43 | + | ||
| 44 | + | /// Where a collected artifact's bytes go. | |
| 45 | + | /// | |
| 46 | + | /// A trait rather than a direct upload call so the download loop can hand each | |
| 47 | + | /// tar off and drop it rather than buffering every artifact — | |
| 48 | + | /// `artifact_run_max_mb` defaults to 512, which is not an amount to hold in | |
| 49 | + | /// RAM. Returning [`Stored`] rather than `()` is what lets the caller charge | |
| 50 | + | /// the run budget the size that actually landed on anvil's disk, which the | |
| 51 | + | /// runner cannot compute itself: a `browse` directory is extracted there and | |
| 52 | + | /// any other directory recompressed. | |
| 53 | + | // async-trait rewrites the method to return a boxed future, which is already | |
| 54 | + | // `#[must_use]`; the attribute it also emits then trips `double_must_use`. | |
| 55 | + | #[allow(clippy::double_must_use)] | |
| 56 | + | #[async_trait::async_trait] | |
| 57 | + | pub trait ArtifactSink: Send { | |
| 58 | + | async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String>; | |
| 59 | + | } | |
| 60 | + | ||
| 61 | + | /// Execute a job in a sandboxed container, streaming output into `log`. | |
| 62 | + | /// Returns the exit code and whatever artifacts `sink` accepted (none on | |
| 63 | + | /// timeout — the container is already gone). | |
| 64 | + | /// | |
| 65 | + | /// The job container never sees the Docker socket and gets no mounts of any | |
| 66 | + | /// kind (the checkout is *uploaded*, not bind-mounted; artifacts are | |
| 67 | + | /// *downloaded* out the same way). All capabilities are dropped and | |
| 68 | + | /// `no-new-privileges` is set unconditionally; pids/memory/cpu caps, the | |
| 69 | + | /// wall-clock timeout, network access and the container user come from | |
| 70 | + | /// `spec.sandbox`. The image allowlist was applied when the spec was built — | |
| 71 | + | /// a runner receiving a spec does not get to widen it. | |
| 72 | + | pub async fn execute( | |
| 73 | + | spec: &JobSpec, | |
| 74 | + | tar: Vec<u8>, | |
| 75 | + | log: &mut String, | |
| 76 | + | sink: &mut dyn ArtifactSink, | |
| 77 | + | ) -> Result<(i64, Vec<CollectedArtifact>), String> { | |
| 78 | + | let platform = spec.platform.clone().unwrap_or_default(); | |
| 79 | + | let docker = anvil_docker::connect()?; | |
| 80 | + | anvil_docker::ensure_image(&docker, &spec.image, &platform).await?; | |
| 81 | + | let sb = &spec.sandbox; | |
| 82 | + | ||
| 83 | + | // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap. | |
| 84 | + | let host_config = HostConfig { | |
| 85 | + | cap_drop: Some(vec!["ALL".to_string()]), | |
| 86 | + | security_opt: Some(vec!["no-new-privileges:true".to_string()]), | |
| 87 | + | pids_limit: (sb.pids_limit > 0).then_some(sb.pids_limit), | |
| 88 | + | memory: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024), | |
| 89 | + | memory_swap: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024), | |
| 90 | + | nano_cpus: (sb.cpus > 0.0).then_some((sb.cpus * 1e9) as i64), | |
| 91 | + | network_mode: (!sb.network).then(|| "none".to_string()), | |
| 92 | + | ..Default::default() | |
| 93 | + | }; | |
| 94 | + | let config = Config { | |
| 95 | + | image: Some(spec.image.clone()), | |
| 96 | + | cmd: Some(vec![ | |
| 97 | + | "sh".to_string(), | |
| 98 | + | "-c".to_string(), | |
| 99 | + | spec.script.clone(), | |
| 100 | + | ]), | |
| 101 | + | env: (!spec.env.is_empty()) | |
| 102 | + | .then(|| spec.env.iter().map(|(k, v)| format!("{k}={v}")).collect()), | |
| 103 | + | working_dir: Some(WORKDIR.to_string()), | |
| 104 | + | user: (!sb.run_as.is_empty()).then(|| sb.run_as.clone()), | |
| 105 | + | host_config: Some(host_config), | |
| 106 | + | ..Default::default() | |
| 107 | + | }; | |
| 108 | + | // An explicit platform needs the options struct; without one, pass None so | |
| 109 | + | // the daemon picks its native architecture exactly as before. | |
| 110 | + | let opts = spec.platform.as_ref().map(|p| CreateContainerOptions { | |
| 111 | + | name: String::new(), | |
| 112 | + | platform: Some(p.clone()), | |
| 113 | + | }); | |
| 114 | + | let created = docker | |
| 115 | + | .create_container(opts, config) | |
| 116 | + | .await | |
| 117 | + | .map_err(|e| format!("create container: {e}"))?; | |
| 118 | + | let id = created.id; | |
| 119 | + | ||
| 120 | + | // Upload the checkout (tar entries are under `workspace/`, extracted at `/`). | |
| 121 | + | docker | |
| 122 | + | .upload_to_container( | |
| 123 | + | &id, | |
| 124 | + | Some(UploadToContainerOptions { | |
| 125 | + | path: "/".to_string(), | |
| 126 | + | ..Default::default() | |
| 127 | + | }), | |
| 128 | + | tar.into(), | |
| 129 | + | ) | |
| 130 | + | .await | |
| 131 | + | .map_err(|e| format!("upload checkout: {e}"))?; | |
| 132 | + | ||
| 133 | + | docker | |
| 134 | + | .start_container(&id, None::<StartContainerOptions<String>>) | |
| 135 | + | .await | |
| 136 | + | .map_err(|e| format!("start container: {e}"))?; | |
| 137 | + | ||
| 138 | + | // Stream logs and wait for the exit code, bounded by the wall-clock | |
| 139 | + | // timeout. The container is force-removed on every path (which also kills | |
| 140 | + | // a still-running job after a timeout). | |
| 141 | + | let run = async { | |
| 142 | + | let mut logs = docker.logs( | |
| 143 | + | &id, | |
| 144 | + | Some(LogsOptions::<String> { | |
| 145 | + | follow: true, | |
| 146 | + | stdout: true, | |
| 147 | + | stderr: true, | |
| 148 | + | ..Default::default() | |
| 149 | + | }), | |
| 150 | + | ); | |
| 151 | + | while let Some(item) = logs.next().await { | |
| 152 | + | match item { | |
| 153 | + | Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())), | |
| 154 | + | Err(e) => { | |
| 155 | + | log.push_str(&format!("\n[log stream error] {e}\n")); | |
| 156 | + | break; | |
| 157 | + | } | |
| 158 | + | } | |
| 159 | + | } | |
| 160 | + | ||
| 161 | + | // Non-zero exit codes surface as a wait error in bollard. | |
| 162 | + | let mut code = 0i64; | |
| 163 | + | let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>); | |
| 164 | + | while let Some(item) = wait.next().await { | |
| 165 | + | match item { | |
| 166 | + | Ok(resp) => code = resp.status_code, | |
| 167 | + | Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c, | |
| 168 | + | Err(e) => return Err(format!("wait: {e}")), | |
| 169 | + | } | |
| 170 | + | } | |
| 171 | + | Ok(code) | |
| 172 | + | }; | |
| 173 | + | let result = match sb.timeout_secs { | |
| 174 | + | 0 => run.await, | |
| 175 | + | secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run) | |
| 176 | + | .await | |
| 177 | + | .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))), | |
| 178 | + | }; | |
| 179 | + | ||
| 180 | + | // Artifacts come out of the (now stopped) container before it is removed. | |
| 181 | + | let collected = match &result { | |
| 182 | + | Ok(_) if !spec.artifacts.is_empty() => { | |
| 183 | + | collect_artifacts(&docker, &id, spec, log, sink).await | |
| 184 | + | } | |
| 185 | + | _ => Vec::new(), | |
| 186 | + | }; | |
| 187 | + | ||
| 188 | + | let _ = docker | |
| 189 | + | .remove_container( | |
| 190 | + | &id, | |
| 191 | + | Some(RemoveContainerOptions { | |
| 192 | + | force: true, | |
| 193 | + | ..Default::default() | |
| 194 | + | }), | |
| 195 | + | ) | |
| 196 | + | .await; | |
| 197 | + | ||
| 198 | + | result.map(|code| (code, collected)) | |
| 199 | + | } | |
| 200 | + | ||
| 201 | + | /// Collect the job's declared artifacts from the stopped container, handing | |
| 202 | + | /// each tar to `sink` as it is downloaded. Failures are per-artifact: each is | |
| 203 | + | /// logged and skipped, never failing the run. | |
| 204 | + | /// | |
| 205 | + | /// One artifact is in memory at a time by construction — the tar is dropped | |
| 206 | + | /// once the sink has taken it — which is what keeps `artifact_run_max_mb` | |
| 207 | + | /// (512 MiB by default) a disk budget rather than a memory one. | |
| 208 | + | async fn collect_artifacts( | |
| 209 | + | docker: &Docker, | |
| 210 | + | id: &str, | |
| 211 | + | job: &JobSpec, | |
| 212 | + | log: &mut String, | |
| 213 | + | sink: &mut dyn ArtifactSink, | |
| 214 | + | ) -> Vec<CollectedArtifact> { | |
| 215 | + | // Extractor outputs first: artifact name → key → value. | |
| 216 | + | let mut metas: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); | |
| 217 | + | if job.artifacts.iter().any(|a| a.has_meta) { | |
| 218 | + | match download_tar(docker, id, META_DIR, META_TAR_CAP).await { | |
| 219 | + | Ok(Some(bytes)) => metas = parse_meta_tar(&bytes), | |
| 220 | + | Ok(None) => log.push_str("\n[artifacts: extractor output exceeded its cap]\n"), | |
| 221 | + | Err(e) => log.push_str(&format!("\n[artifacts: reading extractor output: {e}]\n")), | |
| 222 | + | } | |
| 223 | + | } | |
| 224 | + | ||
| 225 | + | let per_artifact_cap = mb_cap(job.sandbox.artifact_max_mb); | |
| 226 | + | let mut run_budget = mb_cap(job.sandbox.artifact_run_max_mb); | |
| 227 | + | let mut collected = Vec::new(); | |
| 228 | + | for spec in &job.artifacts { | |
| 229 | + | let cap = per_artifact_cap.min(run_budget); | |
| 230 | + | let note = |log: &mut String, what: &str| { | |
| 231 | + | log.push_str(&format!("\n[artifact {}: {what}]\n", spec.name)); | |
| 232 | + | }; | |
| 233 | + | let bytes = match download_tar(docker, id, &format!("{WORKDIR}/{}", spec.path), cap).await { | |
| 234 | + | Ok(Some(bytes)) => bytes, | |
| 235 | + | Ok(None) => { | |
| 236 | + | note(log, "exceeds the size cap; skipped"); | |
| 237 | + | continue; | |
| 238 | + | } | |
| 239 | + | Err(e) => { | |
| 240 | + | note(log, &format!("download failed ({e}); skipped")); | |
| 241 | + | continue; | |
| 242 | + | } | |
| 243 | + | }; | |
| 244 | + | match sink.put(spec, &bytes).await { | |
| 245 | + | Ok(Stored { size, is_dir }) => { | |
| 246 | + | run_budget = run_budget.saturating_sub(size as u64); | |
| 247 | + | let meta = metas.get(&spec.name).cloned().unwrap_or_default(); | |
| 248 | + | collected.push(CollectedArtifact { | |
| 249 | + | name: spec.name.clone(), | |
| 250 | + | size, | |
| 251 | + | is_dir, | |
| 252 | + | browse: spec.browse, | |
| 253 | + | meta: serde_json::to_string(&meta).unwrap_or_else(|_| "{}".into()), | |
| 254 | + | }); | |
| 255 | + | } | |
| 256 | + | Err(e) => note(log, &format!("storing failed ({e}); skipped")), | |
| 257 | + | } | |
| 258 | + | } | |
| 259 | + | collected | |
| 260 | + | } | |
| 261 | + | ||
| 262 | + | /// Download `path` from the container as a tar, buffering at most `cap` bytes | |
| 263 | + | /// (`Ok(None)` when exceeded). | |
| 264 | + | async fn download_tar( | |
| 265 | + | docker: &Docker, | |
| 266 | + | id: &str, | |
| 267 | + | path: &str, | |
| 268 | + | cap: u64, | |
| 269 | + | ) -> Result<Option<Vec<u8>>, String> { | |
| 270 | + | let mut stream = docker.download_from_container( | |
| 271 | + | id, | |
| 272 | + | Some(DownloadFromContainerOptions { | |
| 273 | + | path: path.to_string(), | |
| 274 | + | }), | |
| 275 | + | ); | |
| 276 | + | let mut buf = Vec::new(); | |
| 277 | + | while let Some(chunk) = stream.next().await { | |
| 278 | + | let chunk = chunk.map_err(|e| e.to_string())?; | |
| 279 | + | if (buf.len() + chunk.len()) as u64 > cap { | |
| 280 | + | return Ok(None); | |
| 281 | + | } | |
| 282 | + | buf.extend_from_slice(&chunk); | |
| 283 | + | } | |
| 284 | + | Ok(Some(buf)) | |
| 285 | + | } | |
| 286 | + | ||
| 287 | + | /// Parse the extractor-output tar (`anvil-meta/<artifact>/<key>` files) into | |
| 288 | + | /// artifact → key → trimmed value. | |
| 289 | + | fn parse_meta_tar(bytes: &[u8]) -> BTreeMap<String, BTreeMap<String, String>> { | |
| 290 | + | let mut out: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); | |
| 291 | + | let mut archive = tar::Archive::new(bytes); | |
| 292 | + | let Ok(entries) = archive.entries() else { | |
| 293 | + | return out; | |
| 294 | + | }; | |
| 295 | + | for entry in entries.flatten() { | |
| 296 | + | if !entry.header().entry_type().is_file() { | |
| 297 | + | continue; | |
| 298 | + | } | |
| 299 | + | let Ok(path) = entry.path() else { continue }; | |
| 300 | + | // anvil-meta/<artifact>/<key> | |
| 301 | + | let parts: Vec<String> = path | |
| 302 | + | .components() | |
| 303 | + | .skip(1) | |
| 304 | + | .map(|c| c.as_os_str().to_string_lossy().into_owned()) | |
| 305 | + | .collect(); | |
| 306 | + | let [artifact, key] = parts.as_slice() else { | |
| 307 | + | continue; | |
| 308 | + | }; | |
| 309 | + | let (artifact, key) = (artifact.clone(), key.clone()); | |
| 310 | + | let mut value = String::new(); | |
| 311 | + | let _ = entry.take(META_VALUE_CAP as u64).read_to_string(&mut value); | |
| 312 | + | let value = value.trim().to_string(); | |
| 313 | + | if !value.is_empty() { | |
| 314 | + | out.entry(artifact).or_default().insert(key, value); | |
| 315 | + | } | |
| 316 | + | } | |
| 317 | + | out | |
| 318 | + | } | |
| 319 | + | ||
| 320 | + | #[cfg(test)] | |
| 321 | + | mod tests { | |
| 322 | + | use super::*; | |
| 323 | + | ||
| 324 | + | /// A tar with the same layout the Docker archive endpoint produces: | |
| 325 | + | /// entries rooted at the requested item's basename. `None` content means a | |
| 326 | + | /// directory entry. | |
| 327 | + | fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> { | |
| 328 | + | let mut b = tar::Builder::new(Vec::new()); | |
| 329 | + | for (path, content) in entries { | |
| 330 | + | let mut h = tar::Header::new_gnu(); | |
| 331 | + | match content { | |
| 332 | + | Some(c) => { | |
| 333 | + | h.set_size(c.len() as u64); | |
| 334 | + | h.set_mode(0o644); | |
| 335 | + | h.set_entry_type(tar::EntryType::Regular); | |
| 336 | + | b.append_data(&mut h, path, c.as_bytes()).unwrap(); | |
| 337 | + | } | |
| 338 | + | None => { | |
| 339 | + | h.set_size(0); | |
| 340 | + | h.set_mode(0o755); | |
| 341 | + | h.set_entry_type(tar::EntryType::Directory); | |
| 342 | + | b.append_data(&mut h, path, std::io::empty()).unwrap(); | |
| 343 | + | } | |
| 344 | + | } | |
| 345 | + | } | |
| 346 | + | b.into_inner().unwrap() | |
| 347 | + | } | |
| 348 | + | ||
| 349 | + | #[test] | |
| 350 | + | fn parses_meta_tar_with_trimmed_capped_values() { | |
| 351 | + | let tar = tar_of(&[ | |
| 352 | + | ("anvil-meta", None), | |
| 353 | + | ("anvil-meta/bin", None), | |
| 354 | + | ("anvil-meta/bin/version", Some("anvild 0.0.0\n")), | |
| 355 | + | ("anvil-meta/bin/empty", Some(" \n")), | |
| 356 | + | ("anvil-meta/doc", None), | |
| 357 | + | ("anvil-meta/doc/pages", Some("42")), | |
| 358 | + | ]); | |
| 359 | + | let metas = parse_meta_tar(&tar); | |
| 360 | + | assert_eq!(metas["bin"]["version"], "anvild 0.0.0"); | |
| 361 | + | assert_eq!(metas["doc"]["pages"], "42"); | |
| 362 | + | assert!(!metas["bin"].contains_key("empty"), "blank values dropped"); | |
| 363 | + | } | |
| 364 | + | } |
addedcrates/anvil-worker/src/main.rs+205 −0
| 1 | + | //! anvil's job runner. | |
| 2 | + | //! | |
| 3 | + | //! Claims CI jobs from a forge and runs them in sandboxed containers on this | |
| 4 | + | //! machine's Docker daemon. Runs one job at a time, forever, reconnecting | |
| 5 | + | //! through anything that goes wrong. | |
| 6 | + | //! | |
| 7 | + | //! ```sh | |
| 8 | + | //! anvil-worker --url https://anvil.example.com --token "$ANVIL_RUNNER_TOKEN" | |
| 9 | + | //! ``` | |
| 10 | + | //! | |
| 11 | + | //! On macOS run it natively (a launchd agent), not in a container: a | |
| 12 | + | //! containerized runner would need the Docker socket mounted into it, which | |
| 13 | + | //! rebuilds exactly the root-equivalent hole moving execution off the forge | |
| 14 | + | //! was meant to remove. Isolation is unaffected by the host being a Mac — | |
| 15 | + | //! Docker Desktop runs every container inside one Linux VM, so the sandbox is | |
| 16 | + | //! enforced by the same kernel primitives as on Linux, with the VM as an extra | |
| 17 | + | //! boundary a bare-metal Linux host does not have. | |
| 18 | + | ||
| 19 | + | mod client; | |
| 20 | + | mod executor; | |
| 21 | + | ||
| 22 | + | use std::{ | |
| 23 | + | sync::Arc, | |
| 24 | + | time::Duration, | |
| 25 | + | }; | |
| 26 | + | ||
| 27 | + | use anvil_job::{ | |
| 28 | + | JobResult, | |
| 29 | + | JobSpec, | |
| 30 | + | RunnerInfo, | |
| 31 | + | }; | |
| 32 | + | use clap::Parser; | |
| 33 | + | use client::{ | |
| 34 | + | Client, | |
| 35 | + | UploadSink, | |
| 36 | + | }; | |
| 37 | + | ||
| 38 | + | /// How long to wait before retrying after a failed claim. Long enough not to | |
| 39 | + | /// hammer a forge that is down or misconfigured; short enough that a runner | |
| 40 | + | /// picks up again promptly once it is back. | |
| 41 | + | const RETRY_DELAY: Duration = Duration::from_secs(15); | |
| 42 | + | ||
| 43 | + | #[derive(Parser)] | |
| 44 | + | #[command(name = "anvil-worker", about = "Claim and run anvil CI jobs")] | |
| 45 | + | struct Args { | |
| 46 | + | /// Base URL of the anvil instance, e.g. https://anvil.example.com | |
| 47 | + | #[arg(long, env = "ANVIL_URL")] | |
| 48 | + | url: String, | |
| 49 | + | ||
| 50 | + | /// Shared secret matching the forge's `[ci] runner_token`. | |
| 51 | + | #[arg(long, env = "ANVIL_RUNNER_TOKEN", hide_env_values = true)] | |
| 52 | + | token: String, | |
| 53 | + | ||
| 54 | + | /// Name for this runner in logs and run headers. Defaults to the hostname. | |
| 55 | + | #[arg(long, env = "ANVIL_RUNNER_NAME")] | |
| 56 | + | name: Option<String>, | |
| 57 | + | } | |
| 58 | + | ||
| 59 | + | #[tokio::main] | |
| 60 | + | async fn main() { | |
| 61 | + | tracing_subscriber::fmt() | |
| 62 | + | .with_env_filter( | |
| 63 | + | tracing_subscriber::EnvFilter::try_from_default_env() | |
| 64 | + | .unwrap_or_else(|_| "anvil_worker=info".into()), | |
| 65 | + | ) | |
| 66 | + | .init(); | |
| 67 | + | ||
| 68 | + | let args = Args::parse(); | |
| 69 | + | let name = args.name.unwrap_or_else(hostname); | |
| 70 | + | let info = RunnerInfo { | |
| 71 | + | name: name.clone(), | |
| 72 | + | platform: native_platform(), | |
| 73 | + | version: env!("CARGO_PKG_VERSION").to_string(), | |
| 74 | + | }; | |
| 75 | + | tracing::info!( | |
| 76 | + | "anvil-worker {} ({}) → {}", | |
| 77 | + | info.name, | |
| 78 | + | info.platform, | |
| 79 | + | args.url | |
| 80 | + | ); | |
| 81 | + | ||
| 82 | + | // Fail fast on a daemon that is not there, rather than claiming a job and | |
| 83 | + | // immediately erroring it. | |
| 84 | + | if let Err(e) = anvil_docker::connect() { | |
| 85 | + | tracing::error!("{e}"); | |
| 86 | + | std::process::exit(1); | |
| 87 | + | } | |
| 88 | + | ||
| 89 | + | let client = match Client::new(&args.url, &args.token, info) { | |
| 90 | + | Ok(c) => Arc::new(c), | |
| 91 | + | Err(e) => { | |
| 92 | + | tracing::error!("{e}"); | |
| 93 | + | std::process::exit(1); | |
| 94 | + | } | |
| 95 | + | }; | |
| 96 | + | ||
| 97 | + | loop { | |
| 98 | + | match client.claim().await { | |
| 99 | + | Ok(Some(job)) => run(&client, job).await, | |
| 100 | + | // The poll expired with nothing queued. Straight back in. | |
| 101 | + | Ok(None) => {} | |
| 102 | + | Err(e) => { | |
| 103 | + | tracing::error!("{e}"); | |
| 104 | + | tokio::time::sleep(RETRY_DELAY).await; | |
| 105 | + | } | |
| 106 | + | } | |
| 107 | + | } | |
| 108 | + | } | |
| 109 | + | ||
| 110 | + | /// Run one claimed job and report it. | |
| 111 | + | /// | |
| 112 | + | /// Every failure path still reports: a job anvil handed out and never hears | |
| 113 | + | /// about again sits `running` until its lease expires, which is a slow and | |
| 114 | + | /// confusing way to learn that an image name was wrong. | |
| 115 | + | async fn run(client: &Arc<Client>, job: JobSpec) { | |
| 116 | + | let run_id = job.run_id; | |
| 117 | + | tracing::info!("run {run_id}: claimed ({})", job.image); | |
| 118 | + | ||
| 119 | + | let tar = match client.checkout(run_id).await { | |
| 120 | + | Ok(tar) => tar, | |
| 121 | + | Err(e) => { | |
| 122 | + | report( | |
| 123 | + | client, | |
| 124 | + | run_id, | |
| 125 | + | JobResult { | |
| 126 | + | exit_code: 0, | |
| 127 | + | log: String::new(), | |
| 128 | + | artifacts: Vec::new(), | |
| 129 | + | runner_error: Some(e), | |
| 130 | + | }, | |
| 131 | + | ) | |
| 132 | + | .await; | |
| 133 | + | return; | |
| 134 | + | } | |
| 135 | + | }; | |
| 136 | + | ||
| 137 | + | // Heartbeat for as long as the job runs; aborted below once it is done. | |
| 138 | + | let beat = tokio::spawn(heartbeat(Arc::clone(client), run_id)); | |
| 139 | + | ||
| 140 | + | let mut log = String::new(); | |
| 141 | + | let mut sink = UploadSink { client, run_id }; | |
| 142 | + | let result = match executor::execute(&job, tar, &mut log, &mut sink).await { | |
| 143 | + | Ok((exit_code, artifacts)) => JobResult { | |
| 144 | + | exit_code, | |
| 145 | + | log, | |
| 146 | + | artifacts, | |
| 147 | + | runner_error: None, | |
| 148 | + | }, | |
| 149 | + | Err(e) => JobResult { | |
| 150 | + | exit_code: 0, | |
| 151 | + | log, | |
| 152 | + | artifacts: Vec::new(), | |
| 153 | + | runner_error: Some(e), | |
| 154 | + | }, | |
| 155 | + | }; | |
| 156 | + | beat.abort(); | |
| 157 | + | ||
| 158 | + | tracing::info!( | |
| 159 | + | "run {run_id}: finished (exit {}{})", | |
| 160 | + | result.exit_code, | |
| 161 | + | result | |
| 162 | + | .runner_error | |
| 163 | + | .as_ref() | |
| 164 | + | .map(|e| format!(", runner error: {e}")) | |
| 165 | + | .unwrap_or_default() | |
| 166 | + | ); | |
| 167 | + | report(client, run_id, result).await; | |
| 168 | + | } | |
| 169 | + | ||
| 170 | + | async fn report(client: &Client, run_id: i64, result: JobResult) { | |
| 171 | + | if let Err(e) = client.report(run_id, &result).await { | |
| 172 | + | tracing::error!("run {run_id}: reporting failed: {e}"); | |
| 173 | + | } | |
| 174 | + | } | |
| 175 | + | ||
| 176 | + | /// Keep the lease alive while a job runs. | |
| 177 | + | async fn heartbeat(client: Arc<Client>, run_id: i64) { | |
| 178 | + | let mut tick = tokio::time::interval(anvil_job::HEARTBEAT_INTERVAL); | |
| 179 | + | loop { | |
| 180 | + | tick.tick().await; | |
| 181 | + | if !client.heartbeat(run_id).await { | |
| 182 | + | tracing::warn!("run {run_id}: anvil no longer thinks we hold this job"); | |
| 183 | + | } | |
| 184 | + | } | |
| 185 | + | } | |
| 186 | + | ||
| 187 | + | /// `os/arch` in Docker's spelling, advertised so anvil can eventually route a | |
| 188 | + | /// job that needs a particular architecture to a runner that has it natively. | |
| 189 | + | fn native_platform() -> String { | |
| 190 | + | let arch = match std::env::consts::ARCH { | |
| 191 | + | "x86_64" => "amd64", | |
| 192 | + | "aarch64" => "arm64", | |
| 193 | + | other => other, | |
| 194 | + | }; | |
| 195 | + | // Always `linux`: containers run in a Linux VM on macOS, so the daemon's | |
| 196 | + | // platform is linux there too, whatever the host OS is. | |
| 197 | + | format!("linux/{arch}") | |
| 198 | + | } | |
| 199 | + | ||
| 200 | + | fn hostname() -> String { | |
| 201 | + | std::env::var("HOSTNAME") | |
| 202 | + | .ok() | |
| 203 | + | .filter(|h| !h.is_empty()) | |
| 204 | + | .unwrap_or_else(|| "runner".to_string()) | |
| 205 | + | } |