| 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 | Ok(JobSpec { |
| 69 | run_id, |
| 70 | image: image.to_string(), |
| 71 | // No source for this yet: the pipeline schema has no `platform:` key, |
| 72 | // so every job takes the runner daemon's native architecture, exactly |
| 73 | // as before. Honoured end to end the moment one is set. |
| 74 | platform: None, |
| 75 | script: build_script(pipeline), |
| 76 | env: env.to_vec(), |
| 77 | artifacts: pipeline |
| 78 | .artifacts |
| 79 | .iter() |
| 80 | .map(|a| ArtifactSpec { |
| 81 | name: a.name.clone(), |
| 82 | path: a.path.clone(), |
| 83 | browse: a.browse, |
| 84 | // The extractor commands themselves are already in the script; |
| 85 | // the runner only needs to know whether to look for output. |
| 86 | has_meta: !a.meta.is_empty(), |
| 87 | }) |
| 88 | .collect(), |
| 89 | sandbox: Sandbox { |
| 90 | memory_mb: cfg.memory_mb, |
| 91 | cpus: cfg.cpus, |
| 92 | pids_limit: cfg.pids_limit, |
| 93 | timeout_secs: cfg.timeout_secs, |
| 94 | network: cfg.network, |
| 95 | run_as: cfg.run_as.clone(), |
| 96 | artifact_max_mb: cfg.artifact_max_mb, |
| 97 | artifact_run_max_mb: cfg.artifact_run_max_mb, |
| 98 | }, |
| 99 | }) |
| 100 | } |
| 101 | |
| 102 | /// Assemble the single `sh -c` program a job runs. |
| 103 | /// |
| 104 | /// The steps run in a subshell so the meta-extractor trailer still runs (and |
| 105 | /// the original exit code is preserved) when a step fails — failure artifacts |
| 106 | /// like test reports are the ones that matter most. |
| 107 | fn build_script(pipeline: &Pipeline) -> String { |
| 108 | let mut script = String::from("(\nset -e\n"); |
| 109 | for step in &pipeline.steps { |
| 110 | script.push_str("printf '\\n=== %s ===\\n' "); |
| 111 | script.push_str(&single_quote(step.label())); |
| 112 | script.push('\n'); |
| 113 | script.push_str(&step.run); |
| 114 | script.push('\n'); |
| 115 | } |
| 116 | script.push_str(")\nanvil_rc=$?\n"); |
| 117 | for a in &pipeline.artifacts { |
| 118 | if a.meta.is_empty() { |
| 119 | continue; |
| 120 | } |
| 121 | // Names and keys are parse-time validated to [A-Za-z0-9._-]+, so they |
| 122 | // interpolate into the script safely. |
| 123 | script.push_str(&format!("mkdir -p {META_DIR}/{}\n", a.name)); |
| 124 | for (key, cmd) in &a.meta { |
| 125 | script.push_str(&format!( |
| 126 | "{{\n{cmd}\n}} > {META_DIR}/{}/{key} 2>/dev/null || :\n", |
| 127 | a.name |
| 128 | )); |
| 129 | } |
| 130 | } |
| 131 | script.push_str("exit $anvil_rc\n"); |
| 132 | script |
| 133 | } |
| 134 | |
| 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>) { |
| 143 | match ci::requeue_interrupted(&app.db).await { |
| 144 | Ok(ids) if !ids.is_empty() => { |
| 145 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) |
| 146 | } |
| 147 | Ok(_) => {} |
| 148 | Err(e) => tracing::error!("ci: requeue failed: {e}"), |
| 149 | } |
| 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; |
| 188 | } |
| 189 | app.jobs.wake(); |
| 190 | } |
| 191 | } |
| 192 | } |
| 193 | |
| 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 | // Held across the whole attempt: listing the queue and marking a run |
| 203 | // running are separate awaits, so without it two runners polling at the |
| 204 | // same moment could both be handed the same job. |
| 205 | let _claiming = app.jobs.claim_guard().await; |
| 206 | |
| 207 | let queued = match ci::queued_ids(&app.db).await { |
| 208 | Ok(ids) => ids, |
| 209 | Err(e) => { |
| 210 | tracing::error!("ci: listing queued runs failed: {e}"); |
| 211 | return None; |
| 212 | } |
| 213 | }; |
| 214 | for run_id in queued { |
| 215 | // Skip anything a concurrent claim already took: `queued_ids` reads |
| 216 | // committed state, and `prepare` marks running. |
| 217 | if app.jobs.holds_any(run_id) { |
| 218 | continue; |
| 219 | } |
| 220 | match prepare(app, run_id, runner).await { |
| 221 | Ok(Some(job)) => return Some(job), |
| 222 | Ok(None) => continue, |
| 223 | Err(e) => { |
| 224 | tracing::error!("ci: preparing run {run_id} failed: {e}"); |
| 225 | fail_run(app, run_id, &format!("\n[dispatch error] {e}\n")).await; |
| 226 | } |
| 227 | } |
| 228 | } |
| 229 | None |
| 230 | } |
| 231 | |
| 232 | /// Everything anvil does before a job leaves the building: resolve the repo, |
| 233 | /// parse the pipeline, open the vault, build the spec, take the lease. |
| 234 | /// |
| 235 | /// `Ok(None)` means the run was finished here and should not be dispatched. |
| 236 | async fn prepare(app: &App, run_id: i64, runner: &str) -> Result<Option<JobSpec>, String> { |
| 237 | let run = ci::get(&app.db, run_id) |
| 238 | .await |
| 239 | .map_err(|e| e.to_string())? |
| 240 | .ok_or("run not found")?; |
| 241 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 242 | .await |
| 243 | .map_err(|e| e.to_string())? |
| 244 | .ok_or("repository not found")?; |
| 245 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 246 | .await |
| 247 | .map_err(|e| e.to_string())? |
| 248 | .ok_or("owner not found")?; |
| 249 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 250 | |
| 251 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 252 | .map_err(|e| e.to_string())? |
| 253 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 254 | let pipeline = |
| 255 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 256 | |
| 257 | let short = &run.commit[..run.commit.len().min(12)]; |
| 258 | let mut log = format!( |
| 259 | "anvil ci · {}/{} · {} @ {short}\nrunner: {runner}\nimage: {}\n", |
| 260 | owner.username, |
| 261 | repo.name, |
| 262 | run.ref_name, |
| 263 | app.config.ci.resolve_image(&pipeline.image) |
| 264 | ); |
| 265 | |
| 266 | // Secrets the pipeline asked for, from the in-memory vault. anvil holds no |
| 267 | // key that opens the stored envelopes, so an unlock must have happened |
| 268 | // (`anvild secret unlock`) or the run cannot proceed. |
| 269 | let env = match app.vault.take(run.repo_id, &pipeline.secrets) { |
| 270 | Ok(values) => values, |
| 271 | Err(missing) => { |
| 272 | log.push_str(&format!( |
| 273 | "\n[secrets unavailable: {}]\n\ |
| 274 | This repository is sealed or was unlocked without them. Run:\n \ |
| 275 | anvild secret unlock {}/{}\n", |
| 276 | missing.join(", "), |
| 277 | owner.username, |
| 278 | repo.name, |
| 279 | )); |
| 280 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 281 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 282 | tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name); |
| 283 | return Ok(None); |
| 284 | } |
| 285 | }; |
| 286 | if !env.is_empty() { |
| 287 | log.push_str(&format!( |
| 288 | "secrets: {}\n", |
| 289 | env.iter() |
| 290 | .map(|(name, _)| name.as_str()) |
| 291 | .collect::<Vec<_>>() |
| 292 | .join(", ") |
| 293 | )); |
| 294 | } |
| 295 | |
| 296 | let job = match build_job(run_id, &pipeline, &app.config.ci, &env) { |
| 297 | Ok(job) => job, |
| 298 | Err(e) => { |
| 299 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 300 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 301 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 302 | return Ok(None); |
| 303 | } |
| 304 | }; |
| 305 | |
| 306 | // The header is written now rather than kept in memory until the result |
| 307 | // lands, so a run that is visibly `running` has something to show. |
| 308 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 309 | ci::mark_running(&app.db, run_id).await.ok(); |
| 310 | app.jobs.claim(run_id, runner, env); |
| 311 | tracing::info!("ci: run {run_id} claimed by {runner}"); |
| 312 | Ok(Some(job)) |
| 313 | } |
| 314 | |
| 315 | /// The checkout a claimed job runs against, as an uncompressed tar. |
| 316 | /// |
| 317 | /// Rebuilt from the commit on demand rather than stashed at claim time: it is |
| 318 | /// a pure function of the commit, so a runner retrying the fetch is free, and |
| 319 | /// anvil holds no per-job buffer on a host where memory is the scarce thing. |
| 320 | pub async fn checkout_tar(app: &App, run_id: i64) -> Result<Vec<u8>, String> { |
| 321 | let (_, _, repo_path, run) = resolve(app, run_id).await?; |
| 322 | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; |
| 323 | Ok(build_tar(&files)) |
| 324 | } |
| 325 | |
| 326 | /// The declared spec for one artifact of a run. |
| 327 | /// |
| 328 | /// Re-derived from the pipeline rather than trusted from the upload request: |
| 329 | /// `browse` decides whether a tarball is extracted into a servable directory |
| 330 | /// tree, which is not a choice a runner should get to make. |
| 331 | pub async fn artifact_spec( |
| 332 | app: &App, |
| 333 | run_id: i64, |
| 334 | name: &str, |
| 335 | ) -> Result<Option<ArtifactSpec>, String> { |
| 336 | let (_, _, repo_path, run) = resolve(app, run_id).await?; |
| 337 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 338 | .map_err(|e| e.to_string())? |
| 339 | .ok_or("pipeline missing")?; |
| 340 | let pipeline = |
| 341 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 342 | Ok(pipeline |
| 343 | .artifacts |
| 344 | .iter() |
| 345 | .find(|a| a.name == name) |
| 346 | .map(|a| ArtifactSpec { |
| 347 | name: a.name.clone(), |
| 348 | path: a.path.clone(), |
| 349 | browse: a.browse, |
| 350 | has_meta: !a.meta.is_empty(), |
| 351 | })) |
| 352 | } |
| 353 | |
| 354 | /// Store one uploaded artifact tar into the run's scratch directory. |
| 355 | /// |
| 356 | /// Scratch sits next to the artifacts' final home (same filesystem, so |
| 357 | /// [`finish_run`]'s swap is a rename) and is only moved into place once the |
| 358 | /// run reports in. |
| 359 | pub fn store_upload( |
| 360 | app: &App, |
| 361 | run: &anvil_core::CiRun, |
| 362 | spec: &ArtifactSpec, |
| 363 | tar: &[u8], |
| 364 | ) -> Result<Stored, String> { |
| 365 | let scratch = scratch_dir(app, run); |
| 366 | std::fs::create_dir_all(&scratch).map_err(|e| format!("creating scratch dir: {e}"))?; |
| 367 | let (size, is_dir) = store_artifact(spec, tar, &scratch)?; |
| 368 | Ok(Stored { size, is_dir }) |
| 369 | } |
| 370 | |
| 371 | /// Where a run's artifacts accumulate before the swap. |
| 372 | fn scratch_dir(app: &App, run: &anvil_core::CiRun) -> std::path::PathBuf { |
| 373 | app.config |
| 374 | .artifacts_dir() |
| 375 | .join(run.repo_id.to_string()) |
| 376 | .join(format!(".collecting-{}", run.id)) |
| 377 | } |
| 378 | |
| 379 | /// Record a finished job: swap its artifacts into place, mask and store the |
| 380 | /// log, set the status, and fire the deploy webhook if it qualifies. |
| 381 | pub async fn finish_run(app: &App, run_id: i64, result: JobResult) -> Result<(), String> { |
| 382 | let (owner, repo, repo_path, run) = resolve(app, run_id).await?; |
| 383 | let scratch = scratch_dir(app, &run); |
| 384 | let mut log = result.log; |
| 385 | |
| 386 | let (status, collected) = match result.runner_error { |
| 387 | Some(e) => { |
| 388 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 389 | (ci::status::ERROR, Vec::new()) |
| 390 | } |
| 391 | None if result.exit_code == 0 => (ci::status::SUCCESS, result.artifacts), |
| 392 | None => { |
| 393 | log.push_str(&format!("\n[exited with status {}]\n", result.exit_code)); |
| 394 | (ci::status::FAILURE, result.artifacts) |
| 395 | } |
| 396 | }; |
| 397 | |
| 398 | // Swap the collected set into place, replacing any earlier run's |
| 399 | // artifacts for this commit, then record the rows. |
| 400 | if !collected.is_empty() { |
| 401 | let artifacts_root = app.config.artifacts_dir(); |
| 402 | let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit); |
| 403 | let swap = async { |
| 404 | ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit) |
| 405 | .await |
| 406 | .map_err(|e| e.to_string())?; |
| 407 | if final_dir.exists() { |
| 408 | std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?; |
| 409 | } |
| 410 | std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?; |
| 411 | for c in &collected { |
| 412 | ci::add_artifact( |
| 413 | &app.db, |
| 414 | run_id, |
| 415 | run.repo_id, |
| 416 | &run.commit, |
| 417 | &c.name, |
| 418 | c.size, |
| 419 | c.is_dir, |
| 420 | c.browse, |
| 421 | &c.meta, |
| 422 | ) |
| 423 | .await |
| 424 | .map_err(|e| e.to_string())?; |
| 425 | } |
| 426 | Ok::<_, String>(()) |
| 427 | }; |
| 428 | match swap.await { |
| 429 | Ok(()) => { |
| 430 | log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len())); |
| 431 | gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await; |
| 432 | } |
| 433 | Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")), |
| 434 | } |
| 435 | } |
| 436 | let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away |
| 437 | |
| 438 | // A step that echoes its environment (`set -x`, `curl -v`, a failing |
| 439 | // command that prints its arguments) would otherwise publish the value on |
| 440 | // a page anyone with read access can see. The values come from the lease |
| 441 | // rather than the vault, which may have re-sealed while the job ran. |
| 442 | mask_secrets(&mut log, &app.jobs.secrets_for(run_id)); |
| 443 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 444 | ci::finish(&app.db, run_id, status).await.ok(); |
| 445 | app.jobs.release(run_id); |
| 446 | tracing::info!("ci: run {run_id} {status}"); |
| 447 | |
| 448 | // Continuous deployment: on a green run of the configured deploy repo's |
| 449 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — |
| 450 | // no other repository can trigger it, even with passing CI. |
| 451 | if status == ci::status::SUCCESS && app.config.ci.is_deploy_target(&owner, &repo, &run.ref_name) |
| 452 | { |
| 453 | deploy(app, &owner, &repo, &run).await; |
| 454 | } |
| 455 | Ok(()) |
| 456 | } |
| 457 | |
| 458 | /// Finish a run that never got far enough to produce a result. |
| 459 | async fn fail_run(app: &App, run_id: i64, note: &str) { |
| 460 | ci::append_log(&app.db, run_id, note).await.ok(); |
| 461 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 462 | app.jobs.release(run_id); |
| 463 | } |
| 464 | |
| 465 | /// (owner username, repo name, repo path, run) for a run id. |
| 466 | async fn resolve( |
| 467 | app: &App, |
| 468 | run_id: i64, |
| 469 | ) -> Result<(String, String, std::path::PathBuf, anvil_core::CiRun), String> { |
| 470 | let run = ci::get(&app.db, run_id) |
| 471 | .await |
| 472 | .map_err(|e| e.to_string())? |
| 473 | .ok_or("run not found")?; |
| 474 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 475 | .await |
| 476 | .map_err(|e| e.to_string())? |
| 477 | .ok_or("repository not found")?; |
| 478 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 479 | .await |
| 480 | .map_err(|e| e.to_string())? |
| 481 | .ok_or("owner not found")?; |
| 482 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 483 | Ok((owner.username, repo.name, repo_path, run)) |
| 484 | } |
| 485 | |
| 486 | /// Enforce `[ci] artifact_quota_mb` for one repository: while over budget, |
| 487 | /// delete the oldest commit's artifacts (rows + directory). Branch-tip |
| 488 | /// commits and the just-stored commit are pinned. Deterministic — runs after |
| 489 | /// every artifact-producing run, no background sweeper. Best-effort: failures |
| 490 | /// are logged, never failing the run. |
| 491 | async fn gc_artifacts( |
| 492 | app: &App, |
| 493 | repo_id: i64, |
| 494 | repo_path: &Path, |
| 495 | keep_commit: &str, |
| 496 | log: &mut String, |
| 497 | ) { |
| 498 | let quota = match app.config.ci.artifact_quota_mb { |
| 499 | 0 => return, |
| 500 | mb => mb_cap(mb), |
| 501 | }; |
| 502 | let rows = match ci::artifacts_for_repo(&app.db, repo_id).await { |
| 503 | Ok(rows) => rows, |
| 504 | Err(e) => { |
| 505 | log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n")); |
| 506 | return; |
| 507 | } |
| 508 | }; |
| 509 | |
| 510 | // Per-commit totals and ages. |
| 511 | let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes) |
| 512 | let mut total: u64 = 0; |
| 513 | for row in &rows { |
| 514 | let entry = commits |
| 515 | .entry(row.commit.clone()) |
| 516 | .or_insert((row.created_at, 0)); |
| 517 | entry.0 = entry.0.min(row.created_at); |
| 518 | entry.1 += row.size; |
| 519 | total = total.saturating_add(row.size.max(0) as u64); |
| 520 | } |
| 521 | if total <= quota { |
| 522 | return; |
| 523 | } |
| 524 | |
| 525 | let pinned = browse::branch_tips(repo_path).unwrap_or_default(); |
| 526 | let mut victims: Vec<(i64, String, i64)> = commits |
| 527 | .into_iter() |
| 528 | .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit)) |
| 529 | .map(|(commit, (oldest, bytes))| (oldest, commit, bytes)) |
| 530 | .collect(); |
| 531 | victims.sort(); |
| 532 | |
| 533 | for (_, commit, bytes) in victims { |
| 534 | if total <= quota { |
| 535 | break; |
| 536 | } |
| 537 | let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit); |
| 538 | if let Err(e) = std::fs::remove_dir_all(&dir) { |
| 539 | log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n")); |
| 540 | continue; // keep the rows; retried next run |
| 541 | } |
| 542 | if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await { |
| 543 | log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n")); |
| 544 | continue; |
| 545 | } |
| 546 | total = total.saturating_sub(bytes.max(0) as u64); |
| 547 | log.push_str(&format!( |
| 548 | "\n[artifact gc: dropped {} for old commit {}]\n", |
| 549 | fmt_mb(bytes), |
| 550 | &commit[..commit.len().min(12)] |
| 551 | )); |
| 552 | } |
| 553 | } |
| 554 | |
| 555 | /// Bytes as a short MiB string for log lines. |
| 556 | fn fmt_mb(bytes: i64) -> String { |
| 557 | format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0)) |
| 558 | } |
| 559 | |
| 560 | /// POST the configured deploy webhook. Best-effort: logs success/failure but |
| 561 | /// never fails the run (CI already passed). |
| 562 | async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) { |
| 563 | let cfg = &app.config.ci; |
| 564 | let body = serde_json::json!({ |
| 565 | "repo": format!("{owner}/{name}"), |
| 566 | "ref": run.ref_name, |
| 567 | "commit": run.commit, |
| 568 | "run_id": run.id, |
| 569 | }); |
| 570 | let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body); |
| 571 | if !cfg.deploy_secret.is_empty() { |
| 572 | req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret); |
| 573 | } |
| 574 | match req.send().await { |
| 575 | Ok(resp) if resp.status().is_success() => { |
| 576 | tracing::info!( |
| 577 | "ci: deploy webhook for {owner}/{name} accepted ({})", |
| 578 | resp.status() |
| 579 | ) |
| 580 | } |
| 581 | Ok(resp) => tracing::error!( |
| 582 | "ci: deploy webhook for {owner}/{name} returned {}", |
| 583 | resp.status() |
| 584 | ), |
| 585 | Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"), |
| 586 | } |
| 587 | } |
| 588 | |
| 589 | /// Replace every secret value in `log` with `***`. |
| 590 | /// |
| 591 | /// Only values worth hiding: very short ones (a one-character secret) would |
| 592 | /// mask half the log for no benefit, and are not credentials in practice. |
| 593 | fn mask_secrets(log: &mut String, env: &[(String, String)]) { |
| 594 | for (_, value) in env { |
| 595 | if value.len() >= 4 && log.contains(value.as_str()) { |
| 596 | *log = log.replace(value.as_str(), "***"); |
| 597 | } |
| 598 | } |
| 599 | } |
| 600 | |
| 601 | /// Write one downloaded artifact tar into `scratch`, returning (size, is_dir). |
| 602 | /// |
| 603 | /// - a file is stored as-is at `scratch/<name>` |
| 604 | /// - a directory with `browse` is extracted under `scratch/<name>/` |
| 605 | /// - any other directory is stored compressed at `scratch/<name>.tar.gz` |
| 606 | fn store_artifact( |
| 607 | spec: &ArtifactSpec, |
| 608 | tar_bytes: &[u8], |
| 609 | scratch: &Path, |
| 610 | ) -> Result<(i64, bool), String> { |
| 611 | // The docker archive endpoint roots entries at the requested item's |
| 612 | // basename; its first entry tells file from directory. |
| 613 | let mut archive = tar::Archive::new(tar_bytes); |
| 614 | let mut entries = archive.entries().map_err(|e| e.to_string())?; |
| 615 | let first = entries |
| 616 | .next() |
| 617 | .ok_or("empty archive")? |
| 618 | .map_err(|e| e.to_string())?; |
| 619 | let is_dir = first.header().entry_type().is_dir(); |
| 620 | |
| 621 | if !is_dir { |
| 622 | let mut entry = first; |
| 623 | let mut content = Vec::new(); |
| 624 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 625 | std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?; |
| 626 | return Ok((content.len() as i64, false)); |
| 627 | } |
| 628 | |
| 629 | if !spec.browse { |
| 630 | let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name))) |
| 631 | .map_err(|e| e.to_string())?; |
| 632 | let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default()); |
| 633 | std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?; |
| 634 | let file = enc.finish().map_err(|e| e.to_string())?; |
| 635 | let size = file.metadata().map_err(|e| e.to_string())?.len(); |
| 636 | return Ok((size as i64, true)); |
| 637 | } |
| 638 | |
| 639 | // Browsable: extract regular files under scratch/<name>/, stripping the |
| 640 | // basename prefix. Entry paths come from docker's tar of a real |
| 641 | // filesystem, but stay defensive: relative components only, no links. |
| 642 | let root = scratch.join(&spec.name); |
| 643 | let mut total = 0i64; |
| 644 | let mut archive = tar::Archive::new(tar_bytes); |
| 645 | for entry in archive.entries().map_err(|e| e.to_string())?.flatten() { |
| 646 | if !entry.header().entry_type().is_file() { |
| 647 | continue; |
| 648 | } |
| 649 | let Ok(path) = entry.path() else { continue }; |
| 650 | let mut rel = std::path::PathBuf::new(); |
| 651 | let mut ok = true; |
| 652 | for c in path.components().skip(1) { |
| 653 | match c { |
| 654 | std::path::Component::Normal(p) => rel.push(p), |
| 655 | _ => { |
| 656 | ok = false; |
| 657 | break; |
| 658 | } |
| 659 | } |
| 660 | } |
| 661 | if !ok || rel.as_os_str().is_empty() { |
| 662 | continue; |
| 663 | } |
| 664 | let dest = root.join(&rel); |
| 665 | if let Some(parent) = dest.parent() { |
| 666 | std::fs::create_dir_all(parent).map_err(|e| e.to_string())?; |
| 667 | } |
| 668 | let mut entry = entry; |
| 669 | let mut content = Vec::new(); |
| 670 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 671 | std::fs::write(&dest, &content).map_err(|e| e.to_string())?; |
| 672 | total += content.len() as i64; |
| 673 | } |
| 674 | let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists |
| 675 | Ok((total, true)) |
| 676 | } |
| 677 | |
| 678 | /// Build an uncompressed tar of the checkout, rooted at `workspace/` so it |
| 679 | /// extracts to `/workspace` when uploaded to the container root. |
| 680 | fn build_tar(files: &[TreeFile]) -> Vec<u8> { |
| 681 | let mut builder = tar::Builder::new(Vec::new()); |
| 682 | for f in files { |
| 683 | let mut header = tar::Header::new_gnu(); |
| 684 | header.set_size(f.content.len() as u64); |
| 685 | header.set_mode(if f.executable { 0o755 } else { 0o644 }); |
| 686 | // append_data sets the path and checksum. |
| 687 | let _ = builder.append_data( |
| 688 | &mut header, |
| 689 | format!("workspace/{}", f.path), |
| 690 | f.content.as_slice(), |
| 691 | ); |
| 692 | } |
| 693 | builder.into_inner().unwrap_or_default() |
| 694 | } |
| 695 | |
| 696 | /// Single-quote a string for safe interpolation into a shell command. |
| 697 | fn single_quote(s: &str) -> String { |
| 698 | format!("'{}'", s.replace('\'', "'\\''")) |
| 699 | } |
| 700 | |
| 701 | #[cfg(test)] |
| 702 | mod tests { |
| 703 | use super::*; |
| 704 | |
| 705 | /// Build a tar the way docker's archive endpoint does: entries rooted at |
| 706 | /// the requested item's basename. |
| 707 | fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> { |
| 708 | let mut b = tar::Builder::new(Vec::new()); |
| 709 | for (path, content) in entries { |
| 710 | let mut h = tar::Header::new_gnu(); |
| 711 | match content { |
| 712 | Some(c) => { |
| 713 | h.set_size(c.len() as u64); |
| 714 | h.set_mode(0o644); |
| 715 | h.set_entry_type(tar::EntryType::Regular); |
| 716 | b.append_data(&mut h, path, c.as_bytes()).unwrap(); |
| 717 | } |
| 718 | None => { |
| 719 | h.set_size(0); |
| 720 | h.set_mode(0o755); |
| 721 | h.set_entry_type(tar::EntryType::Directory); |
| 722 | b.append_data(&mut h, path, std::io::empty()).unwrap(); |
| 723 | } |
| 724 | } |
| 725 | } |
| 726 | b.into_inner().unwrap() |
| 727 | } |
| 728 | |
| 729 | fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec { |
| 730 | ArtifactSpec { |
| 731 | name: name.into(), |
| 732 | path: path.into(), |
| 733 | browse, |
| 734 | has_meta: false, |
| 735 | } |
| 736 | } |
| 737 | |
| 738 | #[test] |
| 739 | fn stores_a_file_artifact_as_is() { |
| 740 | let dir = tempfile::tempdir().unwrap(); |
| 741 | let tar = tar_of(&[("anvild", Some("ELF..."))]); |
| 742 | let (size, is_dir) = store_artifact( |
| 743 | &spec("bin", "target/release/anvild", false), |
| 744 | &tar, |
| 745 | dir.path(), |
| 746 | ) |
| 747 | .unwrap(); |
| 748 | assert!(!is_dir); |
| 749 | assert_eq!(size, 6); |
| 750 | assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF..."); |
| 751 | } |
| 752 | |
| 753 | #[test] |
| 754 | fn stores_a_directory_artifact_as_tar_gz() { |
| 755 | let dir = tempfile::tempdir().unwrap(); |
| 756 | let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]); |
| 757 | let (size, is_dir) = |
| 758 | store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap(); |
| 759 | assert!(is_dir); |
| 760 | let stored = dir.path().join("doc.tar.gz"); |
| 761 | assert_eq!(size, stored.metadata().unwrap().len() as i64); |
| 762 | // Round-trips through gzip back to the original tar bytes. |
| 763 | let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap()); |
| 764 | let mut bytes = Vec::new(); |
| 765 | gz.read_to_end(&mut bytes).unwrap(); |
| 766 | assert_eq!(bytes, tar); |
| 767 | } |
| 768 | |
| 769 | #[test] |
| 770 | fn extracts_a_browsable_directory_artifact() { |
| 771 | let dir = tempfile::tempdir().unwrap(); |
| 772 | let tar = tar_of(&[ |
| 773 | ("doc", None), |
| 774 | ("doc/index.html", Some("<html>")), |
| 775 | ("doc/sub", None), |
| 776 | ("doc/sub/page.html", Some("<p>")), |
| 777 | ]); |
| 778 | let (size, is_dir) = |
| 779 | store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap(); |
| 780 | assert!(is_dir); |
| 781 | assert_eq!(size, 6 + 3); |
| 782 | let root = dir.path().join("doc"); |
| 783 | assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>"); |
| 784 | assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>"); |
| 785 | } |
| 786 | } |