| 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. |
| 3 | //! |
| 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`. |
| 8 | //! |
| 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. |
| 15 | |
| 16 | use std::{ |
| 17 | collections::BTreeMap, |
| 18 | io::Read, |
| 19 | path::Path, |
| 20 | }; |
| 21 | |
| 22 | use anvil_core::{ |
| 23 | App, |
| 24 | ci::{ |
| 25 | self, |
| 26 | Pipeline, |
| 27 | }, |
| 28 | config::CiConfig, |
| 29 | repos, |
| 30 | storage, |
| 31 | users, |
| 32 | }; |
| 33 | use anvil_git::browse::{ |
| 34 | self, |
| 35 | TreeFile, |
| 36 | }; |
| 37 | use anvil_job::{ |
| 38 | ArtifactSpec, |
| 39 | JobResult, |
| 40 | JobSpec, |
| 41 | META_DIR, |
| 42 | Sandbox, |
| 43 | Stored, |
| 44 | mb_cap, |
| 45 | }; |
| 46 | use tokio::sync::mpsc::UnboundedReceiver; |
| 47 | |
| 48 | /// Turn a parsed pipeline into the job a runner receives. |
| 49 | /// |
| 50 | /// The server's half of the split (see `docs/remote-runners.md`): image |
| 51 | /// resolution, the allowlist check and script assembly all happen here. A |
| 52 | /// runner therefore never parses `.anvil/ci.yml`, and cannot widen what it was |
| 53 | /// permitted to run by reinterpreting one. |
| 54 | fn build_job( |
| 55 | run_id: i64, |
| 56 | pipeline: &Pipeline, |
| 57 | cfg: &CiConfig, |
| 58 | env: &[(String, String)], |
| 59 | ) -> Result<JobSpec, String> { |
| 60 | // What the pipeline asked for, or the shared runner image when it omitted |
| 61 | // `image:` entirely. |
| 62 | let image = cfg.resolve_image(&pipeline.image); |
| 63 | if !cfg.image_allowed(image) { |
| 64 | return Err(format!( |
| 65 | "image {image} is not permitted by ci.allowed_images" |
| 66 | )); |
| 67 | } |
| 68 | // `platform:`, else `[ci] platform`, else the runner's native one. The |
| 69 | // pipeline's own value was validated at parse time; this catches a |
| 70 | // malformed `[ci] platform`, which nothing else would. |
| 71 | let platform = cfg.resolve_platform(&pipeline.platform); |
| 72 | if let Some(platform) = platform |
| 73 | && !ci::valid_platform(platform) |
| 74 | { |
| 75 | return Err(format!( |
| 76 | "platform {platform} must be os/arch, e.g. linux/amd64 (from [ci] platform)" |
| 77 | )); |
| 78 | } |
| 79 | Ok(JobSpec { |
| 80 | run_id, |
| 81 | image: image.to_string(), |
| 82 | platform: platform.map(str::to_string), |
| 83 | script: build_script(pipeline), |
| 84 | env: env.to_vec(), |
| 85 | artifacts: pipeline |
| 86 | .artifacts |
| 87 | .iter() |
| 88 | .map(|a| ArtifactSpec { |
| 89 | name: a.name.clone(), |
| 90 | path: a.path.clone(), |
| 91 | browse: a.browse, |
| 92 | // The extractor commands themselves are already in the script; |
| 93 | // the runner only needs to know whether to look for output. |
| 94 | has_meta: !a.meta.is_empty(), |
| 95 | }) |
| 96 | .collect(), |
| 97 | sandbox: Sandbox { |
| 98 | memory_mb: cfg.memory_mb, |
| 99 | cpus: cfg.cpus, |
| 100 | pids_limit: cfg.pids_limit, |
| 101 | timeout_secs: cfg.timeout_secs, |
| 102 | network: cfg.network, |
| 103 | run_as: cfg.run_as.clone(), |
| 104 | artifact_max_mb: cfg.artifact_max_mb, |
| 105 | artifact_run_max_mb: cfg.artifact_run_max_mb, |
| 106 | }, |
| 107 | }) |
| 108 | } |
| 109 | |
| 110 | /// Assemble the single `sh -c` program a job runs. |
| 111 | /// |
| 112 | /// The steps run in a subshell so the meta-extractor trailer still runs (and |
| 113 | /// the original exit code is preserved) when a step fails — failure artifacts |
| 114 | /// like test reports are the ones that matter most. |
| 115 | fn build_script(pipeline: &Pipeline) -> String { |
| 116 | let mut script = String::from("(\nset -e\n"); |
| 117 | for step in &pipeline.steps { |
| 118 | script.push_str("printf '\\n=== %s ===\\n' "); |
| 119 | script.push_str(&single_quote(step.label())); |
| 120 | script.push('\n'); |
| 121 | script.push_str(&step.run); |
| 122 | script.push('\n'); |
| 123 | } |
| 124 | script.push_str(")\nanvil_rc=$?\n"); |
| 125 | for a in &pipeline.artifacts { |
| 126 | if a.meta.is_empty() { |
| 127 | continue; |
| 128 | } |
| 129 | // Names and keys are parse-time validated to [A-Za-z0-9._-]+, so they |
| 130 | // interpolate into the script safely. |
| 131 | script.push_str(&format!("mkdir -p {META_DIR}/{}\n", a.name)); |
| 132 | for (key, cmd) in &a.meta { |
| 133 | script.push_str(&format!( |
| 134 | "{{\n{cmd}\n}} > {META_DIR}/{}/{key} 2>/dev/null || :\n", |
| 135 | a.name |
| 136 | )); |
| 137 | } |
| 138 | } |
| 139 | script.push_str("exit $anvil_rc\n"); |
| 140 | script |
| 141 | } |
| 142 | |
| 143 | /// Dispatch queued runs to runners: recover interrupted runs, then wake parked |
| 144 | /// claim requests as work arrives, and requeue jobs whose runner went away. |
| 145 | /// |
| 146 | /// anvild no longer executes anything (see `docs/remote-runners.md`). This |
| 147 | /// loop used to *be* the runner; now it only decides that a run is ready to be |
| 148 | /// claimed. `rx` still carries newly-enqueued run ids from every push path, so |
| 149 | /// nothing that enqueues had to change. |
| 150 | pub async fn run_dispatcher(app: App, mut rx: UnboundedReceiver<i64>) { |
| 151 | match ci::requeue_interrupted(&app.db).await { |
| 152 | Ok(ids) if !ids.is_empty() => { |
| 153 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) |
| 154 | } |
| 155 | Ok(_) => {} |
| 156 | Err(e) => tracing::error!("ci: requeue failed: {e}"), |
| 157 | } |
| 158 | |
| 159 | tokio::spawn(sweep_expired_leases(app.clone())); |
| 160 | |
| 161 | if app.config.ci.runner_token.is_empty() { |
| 162 | tracing::warn!( |
| 163 | "ci: [ci] runner_token is empty, so no runner can authenticate — \ |
| 164 | queued runs will sit until one is set" |
| 165 | ); |
| 166 | } |
| 167 | tracing::info!("ci dispatcher ready"); |
| 168 | |
| 169 | // Anything already queued is claimable the moment a runner asks; the wake |
| 170 | // is only for runners already parked in a long poll. |
| 171 | app.jobs.wake(); |
| 172 | while rx.recv().await.is_some() { |
| 173 | app.jobs.wake(); |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | /// Requeue runs whose runner stopped heartbeating. |
| 178 | /// |
| 179 | /// The case anvild's own crash never had: with execution in-process, a job |
| 180 | /// died exactly when the process did, and `requeue_interrupted` at startup |
| 181 | /// covered it. A remote runner can vanish while anvil stays up. |
| 182 | /// |
| 183 | /// Note this can double-run a job whose runner is alive but unreachable — the |
| 184 | /// container keeps going, and the requeued run may be claimed elsewhere. CI |
| 185 | /// steps are assumed idempotent, and one job at a time makes it unlikely, but |
| 186 | /// it is the honest failure mode of a lease without fencing. |
| 187 | async fn sweep_expired_leases(app: App) { |
| 188 | let mut tick = tokio::time::interval(anvil_core::jobs::LEASE_TTL / 4); |
| 189 | loop { |
| 190 | tick.tick().await; |
| 191 | for (run_id, runner) in app.jobs.take_expired() { |
| 192 | tracing::warn!("ci: run {run_id} lease expired (runner {runner}); requeueing"); |
| 193 | if let Err(e) = ci::requeue(&app.db, run_id).await { |
| 194 | tracing::error!("ci: requeueing {run_id} failed: {e}"); |
| 195 | continue; |
| 196 | } |
| 197 | app.jobs.wake(); |
| 198 | } |
| 199 | } |
| 200 | } |
| 201 | |
| 202 | /// Hand a queued run to `runner`, or `None` when there is nothing for *this* |
| 203 | /// runner to do. |
| 204 | /// |
| 205 | /// `runner_platform` is what the runner advertised (`linux/arm64`, …) and |
| 206 | /// decides which of the queued runs it is offered — see [`claim_order`]. |
| 207 | /// |
| 208 | /// Runs that cannot be dispatched at all — an unparseable pipeline, a sealed |
| 209 | /// vault, an image the allowlist forbids — are finished here rather than |
| 210 | /// handed out, and the next queued run is tried. That keeps a single bad |
| 211 | /// pipeline from wedging the queue. |
| 212 | pub async fn claim_next(app: &App, runner: &str, runner_platform: &str) -> Option<JobSpec> { |
| 213 | // Held across the whole attempt: listing the queue and marking a run |
| 214 | // running are separate awaits, so without it two runners polling at the |
| 215 | // same moment could both be handed the same job. |
| 216 | let _claiming = app.jobs.claim_guard().await; |
| 217 | |
| 218 | let queued = match ci::queued_ids(&app.db).await { |
| 219 | Ok(ids) => ids, |
| 220 | Err(e) => { |
| 221 | tracing::error!("ci: listing queued runs failed: {e}"); |
| 222 | return None; |
| 223 | } |
| 224 | }; |
| 225 | |
| 226 | // What each queued run wants to run on, oldest first. Reading the pipeline |
| 227 | // is what tells us, so this is a git blob read per queued run — cheap, and |
| 228 | // only on the claim path. |
| 229 | let mut candidates = Vec::new(); |
| 230 | for run_id in queued { |
| 231 | // Skip anything a concurrent claim already took: `queued_ids` reads |
| 232 | // committed state, and `prepare` marks running. |
| 233 | if app.jobs.holds_any(run_id) { |
| 234 | continue; |
| 235 | } |
| 236 | match target_platform(app, run_id).await { |
| 237 | Ok(platform) => candidates.push((run_id, platform)), |
| 238 | Err(e) => { |
| 239 | tracing::error!("ci: preparing run {run_id} failed: {e}"); |
| 240 | fail_run(app, run_id, &format!("\n[dispatch error] {e}\n")).await; |
| 241 | } |
| 242 | } |
| 243 | } |
| 244 | |
| 245 | for run_id in claim_order(&candidates, runner_platform, &app.jobs.platforms()) { |
| 246 | match prepare(app, run_id, runner, runner_platform).await { |
| 247 | Ok(Some(job)) => return Some(job), |
| 248 | Ok(None) => continue, |
| 249 | Err(e) => { |
| 250 | tracing::error!("ci: preparing run {run_id} failed: {e}"); |
| 251 | fail_run(app, run_id, &format!("\n[dispatch error] {e}\n")).await; |
| 252 | } |
| 253 | } |
| 254 | } |
| 255 | None |
| 256 | } |
| 257 | |
| 258 | /// Which of `candidates` this runner should be offered, in order. |
| 259 | /// |
| 260 | /// Two tiers, each in queue order: |
| 261 | /// |
| 262 | /// 1. **Native** — the run named this runner's platform, or named none at all. |
| 263 | /// 2. **Orphaned** — the run named a platform that *no runner present right |
| 264 | /// now* is native to. Someone has to run it, and an emulating runner beats |
| 265 | /// a job that sits queued forever. |
| 266 | /// |
| 267 | /// A run whose platform belongs to another live runner is left alone, which is |
| 268 | /// the whole point: an arm64 Mac stops grabbing the amd64 jobs when an amd64 |
| 269 | /// runner exists to take them, and grabs them the moment it does not. |
| 270 | /// |
| 271 | /// Note what this is not: a promise that the job runs natively. Nothing stops |
| 272 | /// an operator from having exactly one arm64 runner and every pipeline asking |
| 273 | /// for amd64 — that configuration emulates, which is correct, since testing |
| 274 | /// the architecture you ship under Rosetta beats testing one you don't. |
| 275 | fn claim_order( |
| 276 | candidates: &[(i64, Option<String>)], |
| 277 | runner_platform: &str, |
| 278 | live_platforms: &[String], |
| 279 | ) -> Vec<i64> { |
| 280 | let native = |wanted: &Option<String>| match wanted { |
| 281 | None => true, |
| 282 | Some(p) => p == runner_platform, |
| 283 | }; |
| 284 | let orphaned = |wanted: &Option<String>| match wanted { |
| 285 | None => false, |
| 286 | Some(p) => !live_platforms.iter().any(|live| live == p), |
| 287 | }; |
| 288 | candidates |
| 289 | .iter() |
| 290 | .filter(|(_, wanted)| native(wanted)) |
| 291 | .chain( |
| 292 | candidates |
| 293 | .iter() |
| 294 | .filter(|(_, wanted)| !native(wanted) && orphaned(wanted)), |
| 295 | ) |
| 296 | .map(|(run_id, _)| *run_id) |
| 297 | .collect() |
| 298 | } |
| 299 | |
| 300 | /// The platform a queued run needs, for the routing decision — `None` when it |
| 301 | /// names none and can run anywhere. |
| 302 | /// |
| 303 | /// Reads the pipeline and nothing else: no vault, no lease, no `mark_running`. |
| 304 | /// [`prepare`] does that for the run that is actually taken. |
| 305 | async fn target_platform(app: &App, run_id: i64) -> Result<Option<String>, String> { |
| 306 | let (_, _, repo_path, run) = resolve(app, run_id).await?; |
| 307 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 308 | .map_err(|e| e.to_string())? |
| 309 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 310 | let pipeline = |
| 311 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 312 | Ok(app |
| 313 | .config |
| 314 | .ci |
| 315 | .resolve_platform(&pipeline.platform) |
| 316 | .map(str::to_string)) |
| 317 | } |
| 318 | |
| 319 | /// Everything anvil does before a job leaves the building: resolve the repo, |
| 320 | /// parse the pipeline, open the vault, build the spec, take the lease. |
| 321 | /// |
| 322 | /// `Ok(None)` means the run was finished here and should not be dispatched. |
| 323 | async fn prepare( |
| 324 | app: &App, |
| 325 | run_id: i64, |
| 326 | runner: &str, |
| 327 | runner_platform: &str, |
| 328 | ) -> Result<Option<JobSpec>, String> { |
| 329 | let run = ci::get(&app.db, run_id) |
| 330 | .await |
| 331 | .map_err(|e| e.to_string())? |
| 332 | .ok_or("run not found")?; |
| 333 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 334 | .await |
| 335 | .map_err(|e| e.to_string())? |
| 336 | .ok_or("repository not found")?; |
| 337 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 338 | .await |
| 339 | .map_err(|e| e.to_string())? |
| 340 | .ok_or("owner not found")?; |
| 341 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 342 | |
| 343 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 344 | .map_err(|e| e.to_string())? |
| 345 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 346 | let pipeline = |
| 347 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 348 | |
| 349 | let short = &run.commit[..run.commit.len().min(12)]; |
| 350 | let mut log = format!( |
| 351 | "anvil ci · {}/{} · {} @ {short}\nrunner: {runner}\nimage: {}\n", |
| 352 | owner.username, |
| 353 | repo.name, |
| 354 | run.ref_name, |
| 355 | app.config.ci.resolve_image(&pipeline.image) |
| 356 | ); |
| 357 | // Say so when the job is about to be emulated. Without this line a slow |
| 358 | // amd64-on-arm64 build looks like a slow build. |
| 359 | if let Some(platform) = app.config.ci.resolve_platform(&pipeline.platform) { |
| 360 | log.push_str(&format!("platform: {platform}")); |
| 361 | if !runner_platform.is_empty() && runner_platform != platform { |
| 362 | log.push_str(&format!(" (emulated on {runner_platform})")); |
| 363 | } |
| 364 | log.push('\n'); |
| 365 | } |
| 366 | |
| 367 | // Secrets the pipeline asked for, from the in-memory vault. anvil holds no |
| 368 | // key that opens the stored envelopes, so an unlock must have happened |
| 369 | // (`anvild secret unlock`) or the run cannot proceed. |
| 370 | let env = match app.vault.take(run.repo_id, &pipeline.secrets) { |
| 371 | Ok(values) => values, |
| 372 | Err(missing) => { |
| 373 | log.push_str(&format!( |
| 374 | "\n[secrets unavailable: {}]\n\ |
| 375 | This repository is sealed or was unlocked without them. Run:\n \ |
| 376 | anvild secret unlock {}/{}\n", |
| 377 | missing.join(", "), |
| 378 | owner.username, |
| 379 | repo.name, |
| 380 | )); |
| 381 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 382 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 383 | tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name); |
| 384 | return Ok(None); |
| 385 | } |
| 386 | }; |
| 387 | if !env.is_empty() { |
| 388 | log.push_str(&format!( |
| 389 | "secrets: {}\n", |
| 390 | env.iter() |
| 391 | .map(|(name, _)| name.as_str()) |
| 392 | .collect::<Vec<_>>() |
| 393 | .join(", ") |
| 394 | )); |
| 395 | } |
| 396 | |
| 397 | let job = match build_job(run_id, &pipeline, &app.config.ci, &env) { |
| 398 | Ok(job) => job, |
| 399 | Err(e) => { |
| 400 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 401 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 402 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 403 | return Ok(None); |
| 404 | } |
| 405 | }; |
| 406 | |
| 407 | // The header is written now rather than kept in memory until the result |
| 408 | // lands, so a run that is visibly `running` has something to show. |
| 409 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 410 | ci::mark_running(&app.db, run_id).await.ok(); |
| 411 | app.jobs.claim(run_id, runner, env); |
| 412 | tracing::info!("ci: run {run_id} claimed by {runner}"); |
| 413 | Ok(Some(job)) |
| 414 | } |
| 415 | |
| 416 | /// The checkout a claimed job runs against, as an uncompressed tar. |
| 417 | /// |
| 418 | /// Rebuilt from the commit on demand rather than stashed at claim time: it is |
| 419 | /// a pure function of the commit, so a runner retrying the fetch is free, and |
| 420 | /// anvil holds no per-job buffer on a host where memory is the scarce thing. |
| 421 | pub async fn checkout_tar(app: &App, run_id: i64) -> Result<Vec<u8>, String> { |
| 422 | let (_, _, repo_path, run) = resolve(app, run_id).await?; |
| 423 | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; |
| 424 | Ok(build_tar(&files)) |
| 425 | } |
| 426 | |
| 427 | /// The declared spec for one artifact of a run. |
| 428 | /// |
| 429 | /// Re-derived from the pipeline rather than trusted from the upload request: |
| 430 | /// `browse` decides whether a tarball is extracted into a servable directory |
| 431 | /// tree, which is not a choice a runner should get to make. |
| 432 | pub async fn artifact_spec( |
| 433 | app: &App, |
| 434 | run_id: i64, |
| 435 | name: &str, |
| 436 | ) -> Result<Option<ArtifactSpec>, String> { |
| 437 | let (_, _, repo_path, run) = resolve(app, run_id).await?; |
| 438 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 439 | .map_err(|e| e.to_string())? |
| 440 | .ok_or("pipeline missing")?; |
| 441 | let pipeline = |
| 442 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 443 | Ok(pipeline |
| 444 | .artifacts |
| 445 | .iter() |
| 446 | .find(|a| a.name == name) |
| 447 | .map(|a| ArtifactSpec { |
| 448 | name: a.name.clone(), |
| 449 | path: a.path.clone(), |
| 450 | browse: a.browse, |
| 451 | has_meta: !a.meta.is_empty(), |
| 452 | })) |
| 453 | } |
| 454 | |
| 455 | /// Store one uploaded artifact tar into the run's scratch directory. |
| 456 | /// |
| 457 | /// Scratch sits next to the artifacts' final home (same filesystem, so |
| 458 | /// [`finish_run`]'s swap is a rename) and is only moved into place once the |
| 459 | /// run reports in. |
| 460 | pub fn store_upload( |
| 461 | app: &App, |
| 462 | run: &anvil_core::CiRun, |
| 463 | spec: &ArtifactSpec, |
| 464 | tar: &[u8], |
| 465 | ) -> Result<Stored, String> { |
| 466 | let scratch = scratch_dir(app, run); |
| 467 | std::fs::create_dir_all(&scratch).map_err(|e| format!("creating scratch dir: {e}"))?; |
| 468 | let (size, is_dir) = store_artifact(spec, tar, &scratch)?; |
| 469 | Ok(Stored { size, is_dir }) |
| 470 | } |
| 471 | |
| 472 | /// Where a run's artifacts accumulate before the swap. |
| 473 | fn scratch_dir(app: &App, run: &anvil_core::CiRun) -> std::path::PathBuf { |
| 474 | app.config |
| 475 | .artifacts_dir() |
| 476 | .join(run.repo_id.to_string()) |
| 477 | .join(format!(".collecting-{}", run.id)) |
| 478 | } |
| 479 | |
| 480 | /// Record a finished job: swap its artifacts into place, mask and store the |
| 481 | /// log, set the status, and fire the deploy webhook if it qualifies. |
| 482 | pub async fn finish_run(app: &App, run_id: i64, result: JobResult) -> Result<(), String> { |
| 483 | let (owner, repo, repo_path, run) = resolve(app, run_id).await?; |
| 484 | let scratch = scratch_dir(app, &run); |
| 485 | let mut log = result.log; |
| 486 | |
| 487 | let (status, collected) = match result.runner_error { |
| 488 | Some(e) => { |
| 489 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 490 | (ci::status::ERROR, Vec::new()) |
| 491 | } |
| 492 | None if result.exit_code == 0 => (ci::status::SUCCESS, result.artifacts), |
| 493 | None => { |
| 494 | log.push_str(&format!("\n[exited with status {}]\n", result.exit_code)); |
| 495 | (ci::status::FAILURE, result.artifacts) |
| 496 | } |
| 497 | }; |
| 498 | |
| 499 | // Swap the collected set into place, replacing any earlier run's |
| 500 | // artifacts for this commit, then record the rows. |
| 501 | if !collected.is_empty() { |
| 502 | let artifacts_root = app.config.artifacts_dir(); |
| 503 | let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit); |
| 504 | let swap = async { |
| 505 | ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit) |
| 506 | .await |
| 507 | .map_err(|e| e.to_string())?; |
| 508 | if final_dir.exists() { |
| 509 | std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?; |
| 510 | } |
| 511 | std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?; |
| 512 | for c in &collected { |
| 513 | ci::add_artifact( |
| 514 | &app.db, |
| 515 | run_id, |
| 516 | run.repo_id, |
| 517 | &run.commit, |
| 518 | &c.name, |
| 519 | c.size, |
| 520 | c.is_dir, |
| 521 | c.browse, |
| 522 | &c.meta, |
| 523 | ) |
| 524 | .await |
| 525 | .map_err(|e| e.to_string())?; |
| 526 | } |
| 527 | Ok::<_, String>(()) |
| 528 | }; |
| 529 | match swap.await { |
| 530 | Ok(()) => { |
| 531 | log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len())); |
| 532 | gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await; |
| 533 | } |
| 534 | Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")), |
| 535 | } |
| 536 | } |
| 537 | let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away |
| 538 | |
| 539 | // A step that echoes its environment (`set -x`, `curl -v`, a failing |
| 540 | // command that prints its arguments) would otherwise publish the value on |
| 541 | // a page anyone with read access can see. The values come from the lease |
| 542 | // rather than the vault, which may have re-sealed while the job ran. |
| 543 | mask_secrets(&mut log, &app.jobs.secrets_for(run_id)); |
| 544 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 545 | ci::finish(&app.db, run_id, status).await.ok(); |
| 546 | app.jobs.release(run_id); |
| 547 | tracing::info!("ci: run {run_id} {status}"); |
| 548 | |
| 549 | // Continuous deployment: on a green run of the configured deploy repo's |
| 550 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — |
| 551 | // no other repository can trigger it, even with passing CI. |
| 552 | if status == ci::status::SUCCESS && app.config.ci.is_deploy_target(&owner, &repo, &run.ref_name) |
| 553 | { |
| 554 | deploy(app, &owner, &repo, &run).await; |
| 555 | } |
| 556 | Ok(()) |
| 557 | } |
| 558 | |
| 559 | /// Finish a run that never got far enough to produce a result. |
| 560 | async fn fail_run(app: &App, run_id: i64, note: &str) { |
| 561 | ci::append_log(&app.db, run_id, note).await.ok(); |
| 562 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 563 | app.jobs.release(run_id); |
| 564 | } |
| 565 | |
| 566 | /// (owner username, repo name, repo path, run) for a run id. |
| 567 | async fn resolve( |
| 568 | app: &App, |
| 569 | run_id: i64, |
| 570 | ) -> Result<(String, String, std::path::PathBuf, anvil_core::CiRun), String> { |
| 571 | let run = ci::get(&app.db, run_id) |
| 572 | .await |
| 573 | .map_err(|e| e.to_string())? |
| 574 | .ok_or("run not found")?; |
| 575 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 576 | .await |
| 577 | .map_err(|e| e.to_string())? |
| 578 | .ok_or("repository not found")?; |
| 579 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 580 | .await |
| 581 | .map_err(|e| e.to_string())? |
| 582 | .ok_or("owner not found")?; |
| 583 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 584 | Ok((owner.username, repo.name, repo_path, run)) |
| 585 | } |
| 586 | |
| 587 | /// Enforce `[ci] artifact_quota_mb` for one repository: while over budget, |
| 588 | /// delete the oldest commit's artifacts (rows + directory). Branch-tip |
| 589 | /// commits and the just-stored commit are pinned. Deterministic — runs after |
| 590 | /// every artifact-producing run, no background sweeper. Best-effort: failures |
| 591 | /// are logged, never failing the run. |
| 592 | async fn gc_artifacts( |
| 593 | app: &App, |
| 594 | repo_id: i64, |
| 595 | repo_path: &Path, |
| 596 | keep_commit: &str, |
| 597 | log: &mut String, |
| 598 | ) { |
| 599 | let quota = match app.config.ci.artifact_quota_mb { |
| 600 | 0 => return, |
| 601 | mb => mb_cap(mb), |
| 602 | }; |
| 603 | let rows = match ci::artifacts_for_repo(&app.db, repo_id).await { |
| 604 | Ok(rows) => rows, |
| 605 | Err(e) => { |
| 606 | log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n")); |
| 607 | return; |
| 608 | } |
| 609 | }; |
| 610 | |
| 611 | // Per-commit totals and ages. |
| 612 | let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes) |
| 613 | let mut total: u64 = 0; |
| 614 | for row in &rows { |
| 615 | let entry = commits |
| 616 | .entry(row.commit.clone()) |
| 617 | .or_insert((row.created_at, 0)); |
| 618 | entry.0 = entry.0.min(row.created_at); |
| 619 | entry.1 += row.size; |
| 620 | total = total.saturating_add(row.size.max(0) as u64); |
| 621 | } |
| 622 | if total <= quota { |
| 623 | return; |
| 624 | } |
| 625 | |
| 626 | let pinned = browse::branch_tips(repo_path).unwrap_or_default(); |
| 627 | let mut victims: Vec<(i64, String, i64)> = commits |
| 628 | .into_iter() |
| 629 | .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit)) |
| 630 | .map(|(commit, (oldest, bytes))| (oldest, commit, bytes)) |
| 631 | .collect(); |
| 632 | victims.sort(); |
| 633 | |
| 634 | for (_, commit, bytes) in victims { |
| 635 | if total <= quota { |
| 636 | break; |
| 637 | } |
| 638 | let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit); |
| 639 | if let Err(e) = std::fs::remove_dir_all(&dir) { |
| 640 | log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n")); |
| 641 | continue; // keep the rows; retried next run |
| 642 | } |
| 643 | if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await { |
| 644 | log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n")); |
| 645 | continue; |
| 646 | } |
| 647 | total = total.saturating_sub(bytes.max(0) as u64); |
| 648 | log.push_str(&format!( |
| 649 | "\n[artifact gc: dropped {} for old commit {}]\n", |
| 650 | fmt_mb(bytes), |
| 651 | &commit[..commit.len().min(12)] |
| 652 | )); |
| 653 | } |
| 654 | } |
| 655 | |
| 656 | /// Bytes as a short MiB string for log lines. |
| 657 | fn fmt_mb(bytes: i64) -> String { |
| 658 | format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0)) |
| 659 | } |
| 660 | |
| 661 | /// POST the configured deploy webhook. Best-effort: logs success/failure but |
| 662 | /// never fails the run (CI already passed). |
| 663 | async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) { |
| 664 | let cfg = &app.config.ci; |
| 665 | let body = serde_json::json!({ |
| 666 | "repo": format!("{owner}/{name}"), |
| 667 | "ref": run.ref_name, |
| 668 | "commit": run.commit, |
| 669 | "run_id": run.id, |
| 670 | }); |
| 671 | let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body); |
| 672 | if !cfg.deploy_secret.is_empty() { |
| 673 | req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret); |
| 674 | } |
| 675 | match req.send().await { |
| 676 | Ok(resp) if resp.status().is_success() => { |
| 677 | tracing::info!( |
| 678 | "ci: deploy webhook for {owner}/{name} accepted ({})", |
| 679 | resp.status() |
| 680 | ) |
| 681 | } |
| 682 | Ok(resp) => tracing::error!( |
| 683 | "ci: deploy webhook for {owner}/{name} returned {}", |
| 684 | resp.status() |
| 685 | ), |
| 686 | Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"), |
| 687 | } |
| 688 | } |
| 689 | |
| 690 | /// Replace every secret value in `log` with `***`. |
| 691 | /// |
| 692 | /// Only values worth hiding: very short ones (a one-character secret) would |
| 693 | /// mask half the log for no benefit, and are not credentials in practice. |
| 694 | fn mask_secrets(log: &mut String, env: &[(String, String)]) { |
| 695 | for (_, value) in env { |
| 696 | if value.len() >= 4 && log.contains(value.as_str()) { |
| 697 | *log = log.replace(value.as_str(), "***"); |
| 698 | } |
| 699 | } |
| 700 | } |
| 701 | |
| 702 | /// Write one downloaded artifact tar into `scratch`, returning (size, is_dir). |
| 703 | /// |
| 704 | /// - a file is stored as-is at `scratch/<name>` |
| 705 | /// - a directory with `browse` is extracted under `scratch/<name>/` |
| 706 | /// - any other directory is stored compressed at `scratch/<name>.tar.gz` |
| 707 | fn store_artifact( |
| 708 | spec: &ArtifactSpec, |
| 709 | tar_bytes: &[u8], |
| 710 | scratch: &Path, |
| 711 | ) -> Result<(i64, bool), String> { |
| 712 | // The docker archive endpoint roots entries at the requested item's |
| 713 | // basename; its first entry tells file from directory. |
| 714 | let mut archive = tar::Archive::new(tar_bytes); |
| 715 | let mut entries = archive.entries().map_err(|e| e.to_string())?; |
| 716 | let first = entries |
| 717 | .next() |
| 718 | .ok_or("empty archive")? |
| 719 | .map_err(|e| e.to_string())?; |
| 720 | let is_dir = first.header().entry_type().is_dir(); |
| 721 | |
| 722 | if !is_dir { |
| 723 | let mut entry = first; |
| 724 | let mut content = Vec::new(); |
| 725 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 726 | std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?; |
| 727 | return Ok((content.len() as i64, false)); |
| 728 | } |
| 729 | |
| 730 | if !spec.browse { |
| 731 | let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name))) |
| 732 | .map_err(|e| e.to_string())?; |
| 733 | let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default()); |
| 734 | std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?; |
| 735 | let file = enc.finish().map_err(|e| e.to_string())?; |
| 736 | let size = file.metadata().map_err(|e| e.to_string())?.len(); |
| 737 | return Ok((size as i64, true)); |
| 738 | } |
| 739 | |
| 740 | // Browsable: extract regular files under scratch/<name>/, stripping the |
| 741 | // basename prefix. Entry paths come from docker's tar of a real |
| 742 | // filesystem, but stay defensive: relative components only, no links. |
| 743 | let root = scratch.join(&spec.name); |
| 744 | let mut total = 0i64; |
| 745 | let mut archive = tar::Archive::new(tar_bytes); |
| 746 | for entry in archive.entries().map_err(|e| e.to_string())?.flatten() { |
| 747 | if !entry.header().entry_type().is_file() { |
| 748 | continue; |
| 749 | } |
| 750 | let Ok(path) = entry.path() else { continue }; |
| 751 | let mut rel = std::path::PathBuf::new(); |
| 752 | let mut ok = true; |
| 753 | for c in path.components().skip(1) { |
| 754 | match c { |
| 755 | std::path::Component::Normal(p) => rel.push(p), |
| 756 | _ => { |
| 757 | ok = false; |
| 758 | break; |
| 759 | } |
| 760 | } |
| 761 | } |
| 762 | if !ok || rel.as_os_str().is_empty() { |
| 763 | continue; |
| 764 | } |
| 765 | let dest = root.join(&rel); |
| 766 | if let Some(parent) = dest.parent() { |
| 767 | std::fs::create_dir_all(parent).map_err(|e| e.to_string())?; |
| 768 | } |
| 769 | let mut entry = entry; |
| 770 | let mut content = Vec::new(); |
| 771 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 772 | std::fs::write(&dest, &content).map_err(|e| e.to_string())?; |
| 773 | total += content.len() as i64; |
| 774 | } |
| 775 | let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists |
| 776 | Ok((total, true)) |
| 777 | } |
| 778 | |
| 779 | /// Build an uncompressed tar of the checkout, rooted at `workspace/` so it |
| 780 | /// extracts to `/workspace` when uploaded to the container root. |
| 781 | fn build_tar(files: &[TreeFile]) -> Vec<u8> { |
| 782 | let mut builder = tar::Builder::new(Vec::new()); |
| 783 | for f in files { |
| 784 | let mut header = tar::Header::new_gnu(); |
| 785 | header.set_size(f.content.len() as u64); |
| 786 | header.set_mode(if f.executable { 0o755 } else { 0o644 }); |
| 787 | // append_data sets the path and checksum. |
| 788 | let _ = builder.append_data( |
| 789 | &mut header, |
| 790 | format!("workspace/{}", f.path), |
| 791 | f.content.as_slice(), |
| 792 | ); |
| 793 | } |
| 794 | builder.into_inner().unwrap_or_default() |
| 795 | } |
| 796 | |
| 797 | /// Single-quote a string for safe interpolation into a shell command. |
| 798 | fn single_quote(s: &str) -> String { |
| 799 | format!("'{}'", s.replace('\'', "'\\''")) |
| 800 | } |
| 801 | |
| 802 | #[cfg(test)] |
| 803 | mod tests { |
| 804 | use super::*; |
| 805 | |
| 806 | /// Build a tar the way docker's archive endpoint does: entries rooted at |
| 807 | /// the requested item's basename. |
| 808 | fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> { |
| 809 | let mut b = tar::Builder::new(Vec::new()); |
| 810 | for (path, content) in entries { |
| 811 | let mut h = tar::Header::new_gnu(); |
| 812 | match content { |
| 813 | Some(c) => { |
| 814 | h.set_size(c.len() as u64); |
| 815 | h.set_mode(0o644); |
| 816 | h.set_entry_type(tar::EntryType::Regular); |
| 817 | b.append_data(&mut h, path, c.as_bytes()).unwrap(); |
| 818 | } |
| 819 | None => { |
| 820 | h.set_size(0); |
| 821 | h.set_mode(0o755); |
| 822 | h.set_entry_type(tar::EntryType::Directory); |
| 823 | b.append_data(&mut h, path, std::io::empty()).unwrap(); |
| 824 | } |
| 825 | } |
| 826 | } |
| 827 | b.into_inner().unwrap() |
| 828 | } |
| 829 | |
| 830 | /// Queued runs as (id, wanted platform), oldest first. |
| 831 | fn queue(items: &[(i64, Option<&str>)]) -> Vec<(i64, Option<String>)> { |
| 832 | items |
| 833 | .iter() |
| 834 | .map(|(id, p)| (*id, p.map(str::to_string))) |
| 835 | .collect() |
| 836 | } |
| 837 | |
| 838 | /// A runner is offered its own platform's jobs and the unconstrained ones, |
| 839 | /// in queue order, ahead of anything it would have to emulate. |
| 840 | #[test] |
| 841 | fn claim_order_prefers_native_work() { |
| 842 | let queued = queue(&[ |
| 843 | (1, Some("linux/amd64")), |
| 844 | (2, None), |
| 845 | (3, Some("linux/arm64")), |
| 846 | (4, Some("linux/amd64")), |
| 847 | ]); |
| 848 | let live = ["linux/amd64".to_string(), "linux/arm64".to_string()]; |
| 849 | |
| 850 | // Both architectures are present, so neither runner touches the |
| 851 | // other's jobs — not even when its own queue is empty. |
| 852 | assert_eq!(claim_order(&queued, "linux/arm64", &live), vec![2, 3]); |
| 853 | assert_eq!(claim_order(&queued, "linux/amd64", &live), vec![1, 2, 4]); |
| 854 | } |
| 855 | |
| 856 | /// With no runner of the named architecture present, the job is offered to |
| 857 | /// whoever asks rather than sitting queued forever — after that runner's |
| 858 | /// own work. This is the single-runner case: one arm64 Mac, pipelines that |
| 859 | /// ask for amd64 because that is what they ship. |
| 860 | #[test] |
| 861 | fn claim_order_falls_back_to_emulation_when_nobody_is_native() { |
| 862 | let queued = queue(&[(1, Some("linux/amd64")), (2, None)]); |
| 863 | let live = ["linux/arm64".to_string()]; |
| 864 | assert_eq!(claim_order(&queued, "linux/arm64", &live), vec![2, 1]); |
| 865 | } |
| 866 | |
| 867 | /// An unknown platform (`linux/riscv64`, a typo) behaves like any other |
| 868 | /// platform nobody is native to: it runs somewhere and fails visibly there, |
| 869 | /// rather than disappearing from the queue. |
| 870 | #[test] |
| 871 | fn claim_order_never_strands_a_run() { |
| 872 | let queued = queue(&[(1, Some("linux/riscv64"))]); |
| 873 | assert_eq!( |
| 874 | claim_order(&queued, "linux/arm64", &["linux/arm64".to_string()]), |
| 875 | vec![1] |
| 876 | ); |
| 877 | } |
| 878 | |
| 879 | fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec { |
| 880 | ArtifactSpec { |
| 881 | name: name.into(), |
| 882 | path: path.into(), |
| 883 | browse, |
| 884 | has_meta: false, |
| 885 | } |
| 886 | } |
| 887 | |
| 888 | #[test] |
| 889 | fn stores_a_file_artifact_as_is() { |
| 890 | let dir = tempfile::tempdir().unwrap(); |
| 891 | let tar = tar_of(&[("anvild", Some("ELF..."))]); |
| 892 | let (size, is_dir) = store_artifact( |
| 893 | &spec("bin", "target/release/anvild", false), |
| 894 | &tar, |
| 895 | dir.path(), |
| 896 | ) |
| 897 | .unwrap(); |
| 898 | assert!(!is_dir); |
| 899 | assert_eq!(size, 6); |
| 900 | assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF..."); |
| 901 | } |
| 902 | |
| 903 | #[test] |
| 904 | fn stores_a_directory_artifact_as_tar_gz() { |
| 905 | let dir = tempfile::tempdir().unwrap(); |
| 906 | let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]); |
| 907 | let (size, is_dir) = |
| 908 | store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap(); |
| 909 | assert!(is_dir); |
| 910 | let stored = dir.path().join("doc.tar.gz"); |
| 911 | assert_eq!(size, stored.metadata().unwrap().len() as i64); |
| 912 | // Round-trips through gzip back to the original tar bytes. |
| 913 | let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap()); |
| 914 | let mut bytes = Vec::new(); |
| 915 | gz.read_to_end(&mut bytes).unwrap(); |
| 916 | assert_eq!(bytes, tar); |
| 917 | } |
| 918 | |
| 919 | #[test] |
| 920 | fn extracts_a_browsable_directory_artifact() { |
| 921 | let dir = tempfile::tempdir().unwrap(); |
| 922 | let tar = tar_of(&[ |
| 923 | ("doc", None), |
| 924 | ("doc/index.html", Some("<html>")), |
| 925 | ("doc/sub", None), |
| 926 | ("doc/sub/page.html", Some("<p>")), |
| 927 | ]); |
| 928 | let (size, is_dir) = |
| 929 | store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap(); |
| 930 | assert!(is_dir); |
| 931 | assert_eq!(size, 6 + 3); |
| 932 | let root = dir.path().join("doc"); |
| 933 | assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>"); |
| 934 | assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>"); |
| 935 | } |
| 936 | } |