| 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). |
| 4 | //! |
| 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. |
| 10 | //! |
| 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`]). |
| 16 | |
| 17 | use std::{ |
| 18 | collections::BTreeMap, |
| 19 | io::Read, |
| 20 | path::Path, |
| 21 | }; |
| 22 | |
| 23 | use anvil_core::{ |
| 24 | App, |
| 25 | ci::{ |
| 26 | self, |
| 27 | ArtifactSpec, |
| 28 | Pipeline, |
| 29 | }, |
| 30 | config::CiConfig, |
| 31 | repos, |
| 32 | storage, |
| 33 | users, |
| 34 | }; |
| 35 | use anvil_git::browse::{ |
| 36 | self, |
| 37 | TreeFile, |
| 38 | }; |
| 39 | use bollard::{ |
| 40 | Docker, |
| 41 | container::{ |
| 42 | Config, |
| 43 | CreateContainerOptions, |
| 44 | DownloadFromContainerOptions, |
| 45 | LogsOptions, |
| 46 | RemoveContainerOptions, |
| 47 | StartContainerOptions, |
| 48 | UploadToContainerOptions, |
| 49 | WaitContainerOptions, |
| 50 | }, |
| 51 | models::HostConfig, |
| 52 | }; |
| 53 | use futures_util::StreamExt; |
| 54 | use tokio::sync::mpsc::UnboundedReceiver; |
| 55 | |
| 56 | const WORKDIR: &str = "/workspace"; |
| 57 | |
| 58 | /// In-container directory where meta-extractor outputs land, one file per |
| 59 | /// `<artifact>/<key>`. Downloaded as a tar after the run; file-per-value |
| 60 | /// sidesteps quoting/JSON-escaping in shell entirely. |
| 61 | const META_DIR: &str = "/tmp/anvil-meta"; |
| 62 | |
| 63 | /// Cap on the meta-extractor tar (the values are short strings). |
| 64 | const META_TAR_CAP: u64 = 1024 * 1024; |
| 65 | |
| 66 | /// Per-value cap on extractor output, in bytes (after trimming). |
| 67 | const META_VALUE_CAP: usize = 1024; |
| 68 | |
| 69 | /// Run the CI worker loop: recover interrupted runs, drain the queue, then |
| 70 | /// process run ids as they arrive on `rx`. Runs one job at a time. |
| 71 | pub async fn run_worker(app: App, mut rx: UnboundedReceiver<i64>) { |
| 72 | match ci::requeue_interrupted(&app.db).await { |
| 73 | Ok(ids) if !ids.is_empty() => { |
| 74 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) |
| 75 | } |
| 76 | Ok(_) => {} |
| 77 | Err(e) => tracing::error!("ci: requeue failed: {e}"), |
| 78 | } |
| 79 | match ci::queued_ids(&app.db).await { |
| 80 | Ok(ids) => { |
| 81 | for id in ids { |
| 82 | run_one(&app, id).await; |
| 83 | } |
| 84 | } |
| 85 | Err(e) => tracing::error!("ci: listing queued runs failed: {e}"), |
| 86 | } |
| 87 | tracing::info!("ci runner ready"); |
| 88 | while let Some(id) = rx.recv().await { |
| 89 | run_one(&app, id).await; |
| 90 | } |
| 91 | } |
| 92 | |
| 93 | async fn run_one(app: &App, run_id: i64) { |
| 94 | tracing::info!("ci: run {run_id} starting"); |
| 95 | if let Err(e) = process(app, run_id).await { |
| 96 | tracing::error!("ci: run {run_id} errored: {e}"); |
| 97 | let _ = ci::append_log(&app.db, run_id, &format!("\n[runner error] {e}\n")).await; |
| 98 | let _ = ci::finish(&app.db, run_id, ci::status::ERROR).await; |
| 99 | } |
| 100 | } |
| 101 | |
| 102 | async fn process(app: &App, run_id: i64) -> Result<(), String> { |
| 103 | let run = ci::get(&app.db, run_id) |
| 104 | .await |
| 105 | .map_err(|e| e.to_string())? |
| 106 | .ok_or("run not found")?; |
| 107 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 108 | .await |
| 109 | .map_err(|e| e.to_string())? |
| 110 | .ok_or("repository not found")?; |
| 111 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 112 | .await |
| 113 | .map_err(|e| e.to_string())? |
| 114 | .ok_or("owner not found")?; |
| 115 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 116 | |
| 117 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 118 | .map_err(|e| e.to_string())? |
| 119 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 120 | let pipeline = |
| 121 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 122 | |
| 123 | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; |
| 124 | let tar = build_tar(&files); |
| 125 | |
| 126 | ci::mark_running(&app.db, run_id).await.ok(); |
| 127 | |
| 128 | let short = &run.commit[..run.commit.len().min(12)]; |
| 129 | let mut log = format!( |
| 130 | "anvil ci · {}/{} · {} @ {short}\nimage: {}\n", |
| 131 | owner.username, |
| 132 | repo.name, |
| 133 | run.ref_name, |
| 134 | app.config.ci.resolve_image(&pipeline.image) |
| 135 | ); |
| 136 | |
| 137 | // Artifacts are collected into a scratch directory next to their final |
| 138 | // home (same filesystem, so the swap below is a rename), then moved into |
| 139 | // place only after the run finishes. |
| 140 | let artifacts_root = app.config.artifacts_dir(); |
| 141 | let scratch = artifacts_root |
| 142 | .join(run.repo_id.to_string()) |
| 143 | .join(format!(".collecting-{run_id}")); |
| 144 | |
| 145 | // Secrets the pipeline asked for, from the in-memory vault. anvil holds no |
| 146 | // key that opens the stored envelopes, so an unlock must have happened |
| 147 | // (`anvild secret unlock`) or the run cannot proceed. |
| 148 | let env = match app.vault.take(run.repo_id, &pipeline.secrets) { |
| 149 | Ok(values) => values, |
| 150 | Err(missing) => { |
| 151 | log.push_str(&format!( |
| 152 | "\n[secrets unavailable: {}]\n\ |
| 153 | This repository is sealed or was unlocked without them. Run:\n \ |
| 154 | anvild secret unlock {}/{}\n", |
| 155 | missing.join(", "), |
| 156 | owner.username, |
| 157 | repo.name, |
| 158 | )); |
| 159 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 160 | ci::finish(&app.db, run_id, ci::status::ERROR).await.ok(); |
| 161 | tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name); |
| 162 | return Ok(()); |
| 163 | } |
| 164 | }; |
| 165 | if !env.is_empty() { |
| 166 | log.push_str(&format!( |
| 167 | "secrets: {}\n", |
| 168 | env.iter() |
| 169 | .map(|(name, _)| name.as_str()) |
| 170 | .collect::<Vec<_>>() |
| 171 | .join(", ") |
| 172 | )); |
| 173 | } |
| 174 | |
| 175 | let (status, collected) = |
| 176 | match execute(&pipeline, tar, &mut log, &app.config.ci, &scratch, &env).await { |
| 177 | Ok((0, collected)) => (ci::status::SUCCESS, collected), |
| 178 | Ok((code, collected)) => { |
| 179 | log.push_str(&format!("\n[exited with status {code}]\n")); |
| 180 | (ci::status::FAILURE, collected) |
| 181 | } |
| 182 | Err(e) => { |
| 183 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 184 | (ci::status::ERROR, Vec::new()) |
| 185 | } |
| 186 | }; |
| 187 | |
| 188 | // Swap the collected set into place, replacing any earlier run's |
| 189 | // artifacts for this commit, then record the rows. |
| 190 | if !collected.is_empty() { |
| 191 | let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit); |
| 192 | let swap = async { |
| 193 | ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit) |
| 194 | .await |
| 195 | .map_err(|e| e.to_string())?; |
| 196 | if final_dir.exists() { |
| 197 | std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?; |
| 198 | } |
| 199 | std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?; |
| 200 | for c in &collected { |
| 201 | ci::add_artifact( |
| 202 | &app.db, |
| 203 | run_id, |
| 204 | run.repo_id, |
| 205 | &run.commit, |
| 206 | &c.name, |
| 207 | c.size, |
| 208 | c.is_dir, |
| 209 | c.browse, |
| 210 | &c.meta, |
| 211 | ) |
| 212 | .await |
| 213 | .map_err(|e| e.to_string())?; |
| 214 | } |
| 215 | Ok::<_, String>(()) |
| 216 | }; |
| 217 | match swap.await { |
| 218 | Ok(()) => { |
| 219 | log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len())); |
| 220 | gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await; |
| 221 | } |
| 222 | Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")), |
| 223 | } |
| 224 | } |
| 225 | let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away |
| 226 | |
| 227 | // A step that echoes its environment (`set -x`, `curl -v`, a failing |
| 228 | // command that prints its arguments) would otherwise publish the value on |
| 229 | // a page anyone with read access can see. |
| 230 | mask_secrets(&mut log, &env); |
| 231 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 232 | ci::finish(&app.db, run_id, status).await.ok(); |
| 233 | tracing::info!("ci: run {run_id} {status}"); |
| 234 | |
| 235 | // Continuous deployment: on a green run of the configured deploy repo's |
| 236 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — |
| 237 | // no other repository can trigger it, even with passing CI. |
| 238 | if status == ci::status::SUCCESS |
| 239 | && app |
| 240 | .config |
| 241 | .ci |
| 242 | .is_deploy_target(&owner.username, &repo.name, &run.ref_name) |
| 243 | { |
| 244 | deploy(app, &owner.username, &repo.name, &run).await; |
| 245 | } |
| 246 | Ok(()) |
| 247 | } |
| 248 | |
| 249 | /// Enforce `[ci] artifact_quota_mb` for one repository: while over budget, |
| 250 | /// delete the oldest commit's artifacts (rows + directory). Branch-tip |
| 251 | /// commits and the just-stored commit are pinned. Deterministic — runs after |
| 252 | /// every artifact-producing run, no background sweeper. Best-effort: failures |
| 253 | /// are logged, never failing the run. |
| 254 | async fn gc_artifacts( |
| 255 | app: &App, |
| 256 | repo_id: i64, |
| 257 | repo_path: &Path, |
| 258 | keep_commit: &str, |
| 259 | log: &mut String, |
| 260 | ) { |
| 261 | let quota = match app.config.ci.artifact_quota_mb { |
| 262 | 0 => return, |
| 263 | mb => mb_cap(mb), |
| 264 | }; |
| 265 | let rows = match ci::artifacts_for_repo(&app.db, repo_id).await { |
| 266 | Ok(rows) => rows, |
| 267 | Err(e) => { |
| 268 | log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n")); |
| 269 | return; |
| 270 | } |
| 271 | }; |
| 272 | |
| 273 | // Per-commit totals and ages. |
| 274 | let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes) |
| 275 | let mut total: u64 = 0; |
| 276 | for row in &rows { |
| 277 | let entry = commits |
| 278 | .entry(row.commit.clone()) |
| 279 | .or_insert((row.created_at, 0)); |
| 280 | entry.0 = entry.0.min(row.created_at); |
| 281 | entry.1 += row.size; |
| 282 | total = total.saturating_add(row.size.max(0) as u64); |
| 283 | } |
| 284 | if total <= quota { |
| 285 | return; |
| 286 | } |
| 287 | |
| 288 | let pinned = browse::branch_tips(repo_path).unwrap_or_default(); |
| 289 | let mut victims: Vec<(i64, String, i64)> = commits |
| 290 | .into_iter() |
| 291 | .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit)) |
| 292 | .map(|(commit, (oldest, bytes))| (oldest, commit, bytes)) |
| 293 | .collect(); |
| 294 | victims.sort(); |
| 295 | |
| 296 | for (_, commit, bytes) in victims { |
| 297 | if total <= quota { |
| 298 | break; |
| 299 | } |
| 300 | let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit); |
| 301 | if let Err(e) = std::fs::remove_dir_all(&dir) { |
| 302 | log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n")); |
| 303 | continue; // keep the rows; retried next run |
| 304 | } |
| 305 | if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await { |
| 306 | log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n")); |
| 307 | continue; |
| 308 | } |
| 309 | total = total.saturating_sub(bytes.max(0) as u64); |
| 310 | log.push_str(&format!( |
| 311 | "\n[artifact gc: dropped {} for old commit {}]\n", |
| 312 | fmt_mb(bytes), |
| 313 | &commit[..commit.len().min(12)] |
| 314 | )); |
| 315 | } |
| 316 | } |
| 317 | |
| 318 | /// Bytes as a short MiB string for log lines. |
| 319 | fn fmt_mb(bytes: i64) -> String { |
| 320 | format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0)) |
| 321 | } |
| 322 | |
| 323 | /// POST the configured deploy webhook. Best-effort: logs success/failure but |
| 324 | /// never fails the run (CI already passed). |
| 325 | async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) { |
| 326 | let cfg = &app.config.ci; |
| 327 | let body = serde_json::json!({ |
| 328 | "repo": format!("{owner}/{name}"), |
| 329 | "ref": run.ref_name, |
| 330 | "commit": run.commit, |
| 331 | "run_id": run.id, |
| 332 | }); |
| 333 | let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body); |
| 334 | if !cfg.deploy_secret.is_empty() { |
| 335 | req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret); |
| 336 | } |
| 337 | match req.send().await { |
| 338 | Ok(resp) if resp.status().is_success() => { |
| 339 | tracing::info!( |
| 340 | "ci: deploy webhook for {owner}/{name} accepted ({})", |
| 341 | resp.status() |
| 342 | ) |
| 343 | } |
| 344 | Ok(resp) => tracing::error!( |
| 345 | "ci: deploy webhook for {owner}/{name} returned {}", |
| 346 | resp.status() |
| 347 | ), |
| 348 | Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"), |
| 349 | } |
| 350 | } |
| 351 | |
| 352 | /// One artifact collected from the job container, already written under the |
| 353 | /// scratch directory; `process` swaps it into the commit's directory. |
| 354 | struct Collected { |
| 355 | name: String, |
| 356 | size: i64, |
| 357 | is_dir: bool, |
| 358 | browse: bool, |
| 359 | /// JSON object of extractor key → output. |
| 360 | meta: String, |
| 361 | } |
| 362 | |
| 363 | /// Execute the pipeline in a sandboxed container, streaming output into `log`. |
| 364 | /// Returns the container's exit code and any artifacts collected into |
| 365 | /// `scratch` (empty on timeout — the container is already gone). |
| 366 | /// |
| 367 | /// The job container never sees the Docker socket and gets no mounts of any |
| 368 | /// kind (the checkout is *uploaded*, not bind-mounted; artifacts are |
| 369 | /// *downloaded* out the same way). All capabilities are dropped and |
| 370 | /// `no-new-privileges` is set unconditionally; pids/memory/cpu caps, the |
| 371 | /// wall-clock timeout, network access, the container user, and the image |
| 372 | /// allowlist come from `cfg`. |
| 373 | /// Replace every secret value in `log` with `***`. |
| 374 | /// |
| 375 | /// Only values worth hiding: very short ones (a one-character secret) would |
| 376 | /// mask half the log for no benefit, and are not credentials in practice. |
| 377 | fn mask_secrets(log: &mut String, env: &[(String, String)]) { |
| 378 | for (_, value) in env { |
| 379 | if value.len() >= 4 && log.contains(value.as_str()) { |
| 380 | *log = log.replace(value.as_str(), "***"); |
| 381 | } |
| 382 | } |
| 383 | } |
| 384 | |
| 385 | async fn execute( |
| 386 | pipeline: &Pipeline, |
| 387 | tar: Vec<u8>, |
| 388 | log: &mut String, |
| 389 | cfg: &CiConfig, |
| 390 | scratch: &Path, |
| 391 | env: &[(String, String)], |
| 392 | ) -> Result<(i64, Vec<Collected>), String> { |
| 393 | // What the pipeline asked for, or the shared runner image when it omitted |
| 394 | // `image:` entirely. |
| 395 | let image = cfg.resolve_image(&pipeline.image); |
| 396 | if !cfg.image_allowed(image) { |
| 397 | return Err(format!( |
| 398 | "image {image} is not permitted by ci.allowed_images" |
| 399 | )); |
| 400 | } |
| 401 | let docker = anvil_docker::connect()?; |
| 402 | anvil_docker::ensure_image(&docker, image).await?; |
| 403 | |
| 404 | // Build a single `set -e` script from the steps. The steps run in a |
| 405 | // subshell so the meta-extractor trailer still runs (and the original |
| 406 | // exit code is preserved) when a step fails — failure artifacts like test |
| 407 | // reports are the ones that matter most. |
| 408 | let mut script = String::from("(\nset -e\n"); |
| 409 | for step in &pipeline.steps { |
| 410 | script.push_str("printf '\\n=== %s ===\\n' "); |
| 411 | script.push_str(&single_quote(step.label())); |
| 412 | script.push('\n'); |
| 413 | script.push_str(&step.run); |
| 414 | script.push('\n'); |
| 415 | } |
| 416 | script.push_str(")\nanvil_rc=$?\n"); |
| 417 | for a in &pipeline.artifacts { |
| 418 | if a.meta.is_empty() { |
| 419 | continue; |
| 420 | } |
| 421 | // Names and keys are parse-time validated to [A-Za-z0-9._-]+, so they |
| 422 | // interpolate into the script safely. |
| 423 | script.push_str(&format!("mkdir -p {META_DIR}/{}\n", a.name)); |
| 424 | for (key, cmd) in &a.meta { |
| 425 | script.push_str(&format!( |
| 426 | "{{\n{cmd}\n}} > {META_DIR}/{}/{key} 2>/dev/null || :\n", |
| 427 | a.name |
| 428 | )); |
| 429 | } |
| 430 | } |
| 431 | script.push_str("exit $anvil_rc\n"); |
| 432 | |
| 433 | // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap. |
| 434 | let host_config = HostConfig { |
| 435 | cap_drop: Some(vec!["ALL".to_string()]), |
| 436 | security_opt: Some(vec!["no-new-privileges:true".to_string()]), |
| 437 | pids_limit: (cfg.pids_limit > 0).then_some(cfg.pids_limit), |
| 438 | memory: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 439 | memory_swap: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 440 | nano_cpus: (cfg.cpus > 0.0).then_some((cfg.cpus * 1e9) as i64), |
| 441 | network_mode: (!cfg.network).then(|| "none".to_string()), |
| 442 | ..Default::default() |
| 443 | }; |
| 444 | let config = Config { |
| 445 | image: Some(image.to_string()), |
| 446 | cmd: Some(vec!["sh".to_string(), "-c".to_string(), script]), |
| 447 | env: (!env.is_empty()).then(|| env.iter().map(|(k, v)| format!("{k}={v}")).collect()), |
| 448 | working_dir: Some(WORKDIR.to_string()), |
| 449 | user: (!cfg.run_as.is_empty()).then(|| cfg.run_as.clone()), |
| 450 | host_config: Some(host_config), |
| 451 | ..Default::default() |
| 452 | }; |
| 453 | let created = docker |
| 454 | .create_container(None::<CreateContainerOptions<String>>, config) |
| 455 | .await |
| 456 | .map_err(|e| format!("create container: {e}"))?; |
| 457 | let id = created.id; |
| 458 | |
| 459 | // Upload the checkout (tar entries are under `workspace/`, extracted at `/`). |
| 460 | docker |
| 461 | .upload_to_container( |
| 462 | &id, |
| 463 | Some(UploadToContainerOptions { |
| 464 | path: "/".to_string(), |
| 465 | ..Default::default() |
| 466 | }), |
| 467 | tar.into(), |
| 468 | ) |
| 469 | .await |
| 470 | .map_err(|e| format!("upload checkout: {e}"))?; |
| 471 | |
| 472 | docker |
| 473 | .start_container(&id, None::<StartContainerOptions<String>>) |
| 474 | .await |
| 475 | .map_err(|e| format!("start container: {e}"))?; |
| 476 | |
| 477 | // Stream logs and wait for the exit code, bounded by the wall-clock |
| 478 | // timeout. The container is force-removed on every path (which also kills |
| 479 | // a still-running job after a timeout). |
| 480 | let run = async { |
| 481 | let mut logs = docker.logs( |
| 482 | &id, |
| 483 | Some(LogsOptions::<String> { |
| 484 | follow: true, |
| 485 | stdout: true, |
| 486 | stderr: true, |
| 487 | ..Default::default() |
| 488 | }), |
| 489 | ); |
| 490 | while let Some(item) = logs.next().await { |
| 491 | match item { |
| 492 | Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())), |
| 493 | Err(e) => { |
| 494 | log.push_str(&format!("\n[log stream error] {e}\n")); |
| 495 | break; |
| 496 | } |
| 497 | } |
| 498 | } |
| 499 | |
| 500 | // Non-zero exit codes surface as a wait error in bollard. |
| 501 | let mut code = 0i64; |
| 502 | let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>); |
| 503 | while let Some(item) = wait.next().await { |
| 504 | match item { |
| 505 | Ok(resp) => code = resp.status_code, |
| 506 | Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c, |
| 507 | Err(e) => return Err(format!("wait: {e}")), |
| 508 | } |
| 509 | } |
| 510 | Ok(code) |
| 511 | }; |
| 512 | let result = match cfg.timeout_secs { |
| 513 | 0 => run.await, |
| 514 | secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run) |
| 515 | .await |
| 516 | .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))), |
| 517 | }; |
| 518 | |
| 519 | // Artifacts come out of the (now stopped) container before it is removed. |
| 520 | let collected = match &result { |
| 521 | Ok(_) if !pipeline.artifacts.is_empty() => { |
| 522 | collect_artifacts(&docker, &id, pipeline, cfg, scratch, log).await |
| 523 | } |
| 524 | _ => Vec::new(), |
| 525 | }; |
| 526 | |
| 527 | let _ = docker |
| 528 | .remove_container( |
| 529 | &id, |
| 530 | Some(RemoveContainerOptions { |
| 531 | force: true, |
| 532 | ..Default::default() |
| 533 | }), |
| 534 | ) |
| 535 | .await; |
| 536 | |
| 537 | result.map(|code| (code, collected)) |
| 538 | } |
| 539 | |
| 540 | /// Collect the pipeline's declared artifacts from the stopped container into |
| 541 | /// `scratch`. Failures are per-artifact: each is logged and skipped, never |
| 542 | /// failing the run. |
| 543 | async fn collect_artifacts( |
| 544 | docker: &Docker, |
| 545 | id: &str, |
| 546 | pipeline: &Pipeline, |
| 547 | cfg: &CiConfig, |
| 548 | scratch: &Path, |
| 549 | log: &mut String, |
| 550 | ) -> Vec<Collected> { |
| 551 | if let Err(e) = std::fs::create_dir_all(scratch) { |
| 552 | log.push_str(&format!( |
| 553 | "\n[artifacts: creating scratch dir failed: {e}]\n" |
| 554 | )); |
| 555 | return Vec::new(); |
| 556 | } |
| 557 | |
| 558 | // Extractor outputs first: artifact name → key → value. |
| 559 | let mut metas: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); |
| 560 | if pipeline.artifacts.iter().any(|a| !a.meta.is_empty()) { |
| 561 | match download_tar(docker, id, META_DIR, META_TAR_CAP).await { |
| 562 | Ok(Some(bytes)) => metas = parse_meta_tar(&bytes), |
| 563 | Ok(None) => log.push_str("\n[artifacts: extractor output exceeded its cap]\n"), |
| 564 | Err(e) => log.push_str(&format!("\n[artifacts: reading extractor output: {e}]\n")), |
| 565 | } |
| 566 | } |
| 567 | |
| 568 | let per_artifact_cap = mb_cap(cfg.artifact_max_mb); |
| 569 | let mut run_budget = mb_cap(cfg.artifact_run_max_mb); |
| 570 | let mut collected = Vec::new(); |
| 571 | for spec in &pipeline.artifacts { |
| 572 | let cap = per_artifact_cap.min(run_budget); |
| 573 | let note = |log: &mut String, what: &str| { |
| 574 | log.push_str(&format!("\n[artifact {}: {what}]\n", spec.name)); |
| 575 | }; |
| 576 | let bytes = match download_tar(docker, id, &format!("{WORKDIR}/{}", spec.path), cap).await { |
| 577 | Ok(Some(bytes)) => bytes, |
| 578 | Ok(None) => { |
| 579 | note(log, "exceeds the size cap; skipped"); |
| 580 | continue; |
| 581 | } |
| 582 | Err(e) => { |
| 583 | note(log, &format!("download failed ({e}); skipped")); |
| 584 | continue; |
| 585 | } |
| 586 | }; |
| 587 | match store_artifact(spec, &bytes, scratch) { |
| 588 | Ok((size, is_dir)) => { |
| 589 | run_budget = run_budget.saturating_sub(size as u64); |
| 590 | let meta = metas.get(&spec.name).cloned().unwrap_or_default(); |
| 591 | collected.push(Collected { |
| 592 | name: spec.name.clone(), |
| 593 | size, |
| 594 | is_dir, |
| 595 | browse: spec.browse, |
| 596 | meta: serde_json::to_string(&meta).unwrap_or_else(|_| "{}".into()), |
| 597 | }); |
| 598 | } |
| 599 | Err(e) => note(log, &format!("storing failed ({e}); skipped")), |
| 600 | } |
| 601 | } |
| 602 | collected |
| 603 | } |
| 604 | |
| 605 | /// `0` (unlimited) → `u64::MAX`, otherwise MiB → bytes. |
| 606 | fn mb_cap(mb: i64) -> u64 { |
| 607 | if mb <= 0 { |
| 608 | u64::MAX |
| 609 | } else { |
| 610 | mb as u64 * 1024 * 1024 |
| 611 | } |
| 612 | } |
| 613 | |
| 614 | /// Download `path` from the container as a tar, buffering at most `cap` bytes |
| 615 | /// (`Ok(None)` when exceeded). |
| 616 | async fn download_tar( |
| 617 | docker: &Docker, |
| 618 | id: &str, |
| 619 | path: &str, |
| 620 | cap: u64, |
| 621 | ) -> Result<Option<Vec<u8>>, String> { |
| 622 | let mut stream = docker.download_from_container( |
| 623 | id, |
| 624 | Some(DownloadFromContainerOptions { |
| 625 | path: path.to_string(), |
| 626 | }), |
| 627 | ); |
| 628 | let mut buf = Vec::new(); |
| 629 | while let Some(chunk) = stream.next().await { |
| 630 | let chunk = chunk.map_err(|e| e.to_string())?; |
| 631 | if (buf.len() + chunk.len()) as u64 > cap { |
| 632 | return Ok(None); |
| 633 | } |
| 634 | buf.extend_from_slice(&chunk); |
| 635 | } |
| 636 | Ok(Some(buf)) |
| 637 | } |
| 638 | |
| 639 | /// Parse the extractor-output tar (`anvil-meta/<artifact>/<key>` files) into |
| 640 | /// artifact → key → trimmed value. |
| 641 | fn parse_meta_tar(bytes: &[u8]) -> BTreeMap<String, BTreeMap<String, String>> { |
| 642 | let mut out: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new(); |
| 643 | let mut archive = tar::Archive::new(bytes); |
| 644 | let Ok(entries) = archive.entries() else { |
| 645 | return out; |
| 646 | }; |
| 647 | for entry in entries.flatten() { |
| 648 | if !entry.header().entry_type().is_file() { |
| 649 | continue; |
| 650 | } |
| 651 | let Ok(path) = entry.path() else { continue }; |
| 652 | // anvil-meta/<artifact>/<key> |
| 653 | let parts: Vec<String> = path |
| 654 | .components() |
| 655 | .skip(1) |
| 656 | .map(|c| c.as_os_str().to_string_lossy().into_owned()) |
| 657 | .collect(); |
| 658 | let [artifact, key] = parts.as_slice() else { |
| 659 | continue; |
| 660 | }; |
| 661 | let (artifact, key) = (artifact.clone(), key.clone()); |
| 662 | let mut value = String::new(); |
| 663 | let _ = entry.take(META_VALUE_CAP as u64).read_to_string(&mut value); |
| 664 | let value = value.trim().to_string(); |
| 665 | if !value.is_empty() { |
| 666 | out.entry(artifact).or_default().insert(key, value); |
| 667 | } |
| 668 | } |
| 669 | out |
| 670 | } |
| 671 | |
| 672 | /// Write one downloaded artifact tar into `scratch`, returning (size, is_dir). |
| 673 | /// |
| 674 | /// - a file is stored as-is at `scratch/<name>` |
| 675 | /// - a directory with `browse` is extracted under `scratch/<name>/` |
| 676 | /// - any other directory is stored compressed at `scratch/<name>.tar.gz` |
| 677 | fn store_artifact( |
| 678 | spec: &ArtifactSpec, |
| 679 | tar_bytes: &[u8], |
| 680 | scratch: &Path, |
| 681 | ) -> Result<(i64, bool), String> { |
| 682 | // The docker archive endpoint roots entries at the requested item's |
| 683 | // basename; its first entry tells file from directory. |
| 684 | let mut archive = tar::Archive::new(tar_bytes); |
| 685 | let mut entries = archive.entries().map_err(|e| e.to_string())?; |
| 686 | let first = entries |
| 687 | .next() |
| 688 | .ok_or("empty archive")? |
| 689 | .map_err(|e| e.to_string())?; |
| 690 | let is_dir = first.header().entry_type().is_dir(); |
| 691 | |
| 692 | if !is_dir { |
| 693 | let mut entry = first; |
| 694 | let mut content = Vec::new(); |
| 695 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 696 | std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?; |
| 697 | return Ok((content.len() as i64, false)); |
| 698 | } |
| 699 | |
| 700 | if !spec.browse { |
| 701 | let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name))) |
| 702 | .map_err(|e| e.to_string())?; |
| 703 | let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default()); |
| 704 | std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?; |
| 705 | let file = enc.finish().map_err(|e| e.to_string())?; |
| 706 | let size = file.metadata().map_err(|e| e.to_string())?.len(); |
| 707 | return Ok((size as i64, true)); |
| 708 | } |
| 709 | |
| 710 | // Browsable: extract regular files under scratch/<name>/, stripping the |
| 711 | // basename prefix. Entry paths come from docker's tar of a real |
| 712 | // filesystem, but stay defensive: relative components only, no links. |
| 713 | let root = scratch.join(&spec.name); |
| 714 | let mut total = 0i64; |
| 715 | let mut archive = tar::Archive::new(tar_bytes); |
| 716 | for entry in archive.entries().map_err(|e| e.to_string())?.flatten() { |
| 717 | if !entry.header().entry_type().is_file() { |
| 718 | continue; |
| 719 | } |
| 720 | let Ok(path) = entry.path() else { continue }; |
| 721 | let mut rel = std::path::PathBuf::new(); |
| 722 | let mut ok = true; |
| 723 | for c in path.components().skip(1) { |
| 724 | match c { |
| 725 | std::path::Component::Normal(p) => rel.push(p), |
| 726 | _ => { |
| 727 | ok = false; |
| 728 | break; |
| 729 | } |
| 730 | } |
| 731 | } |
| 732 | if !ok || rel.as_os_str().is_empty() { |
| 733 | continue; |
| 734 | } |
| 735 | let dest = root.join(&rel); |
| 736 | if let Some(parent) = dest.parent() { |
| 737 | std::fs::create_dir_all(parent).map_err(|e| e.to_string())?; |
| 738 | } |
| 739 | let mut entry = entry; |
| 740 | let mut content = Vec::new(); |
| 741 | entry.read_to_end(&mut content).map_err(|e| e.to_string())?; |
| 742 | std::fs::write(&dest, &content).map_err(|e| e.to_string())?; |
| 743 | total += content.len() as i64; |
| 744 | } |
| 745 | let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists |
| 746 | Ok((total, true)) |
| 747 | } |
| 748 | |
| 749 | /// Build an uncompressed tar of the checkout, rooted at `workspace/` so it |
| 750 | /// extracts to `/workspace` when uploaded to the container root. |
| 751 | fn build_tar(files: &[TreeFile]) -> Vec<u8> { |
| 752 | let mut builder = tar::Builder::new(Vec::new()); |
| 753 | for f in files { |
| 754 | let mut header = tar::Header::new_gnu(); |
| 755 | header.set_size(f.content.len() as u64); |
| 756 | header.set_mode(if f.executable { 0o755 } else { 0o644 }); |
| 757 | // append_data sets the path and checksum. |
| 758 | let _ = builder.append_data( |
| 759 | &mut header, |
| 760 | format!("workspace/{}", f.path), |
| 761 | f.content.as_slice(), |
| 762 | ); |
| 763 | } |
| 764 | builder.into_inner().unwrap_or_default() |
| 765 | } |
| 766 | |
| 767 | /// Single-quote a string for safe interpolation into a shell command. |
| 768 | fn single_quote(s: &str) -> String { |
| 769 | format!("'{}'", s.replace('\'', "'\\''")) |
| 770 | } |
| 771 | |
| 772 | #[cfg(test)] |
| 773 | mod tests { |
| 774 | use super::*; |
| 775 | |
| 776 | /// Build a tar the way docker's archive endpoint does: entries rooted at |
| 777 | /// the requested item's basename. |
| 778 | fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> { |
| 779 | let mut b = tar::Builder::new(Vec::new()); |
| 780 | for (path, content) in entries { |
| 781 | let mut h = tar::Header::new_gnu(); |
| 782 | match content { |
| 783 | Some(c) => { |
| 784 | h.set_size(c.len() as u64); |
| 785 | h.set_mode(0o644); |
| 786 | h.set_entry_type(tar::EntryType::Regular); |
| 787 | b.append_data(&mut h, path, c.as_bytes()).unwrap(); |
| 788 | } |
| 789 | None => { |
| 790 | h.set_size(0); |
| 791 | h.set_mode(0o755); |
| 792 | h.set_entry_type(tar::EntryType::Directory); |
| 793 | b.append_data(&mut h, path, std::io::empty()).unwrap(); |
| 794 | } |
| 795 | } |
| 796 | } |
| 797 | b.into_inner().unwrap() |
| 798 | } |
| 799 | |
| 800 | fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec { |
| 801 | ArtifactSpec { |
| 802 | name: name.into(), |
| 803 | path: path.into(), |
| 804 | browse, |
| 805 | meta: Default::default(), |
| 806 | } |
| 807 | } |
| 808 | |
| 809 | #[test] |
| 810 | fn stores_a_file_artifact_as_is() { |
| 811 | let dir = tempfile::tempdir().unwrap(); |
| 812 | let tar = tar_of(&[("anvild", Some("ELF..."))]); |
| 813 | let (size, is_dir) = store_artifact( |
| 814 | &spec("bin", "target/release/anvild", false), |
| 815 | &tar, |
| 816 | dir.path(), |
| 817 | ) |
| 818 | .unwrap(); |
| 819 | assert!(!is_dir); |
| 820 | assert_eq!(size, 6); |
| 821 | assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF..."); |
| 822 | } |
| 823 | |
| 824 | #[test] |
| 825 | fn stores_a_directory_artifact_as_tar_gz() { |
| 826 | let dir = tempfile::tempdir().unwrap(); |
| 827 | let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]); |
| 828 | let (size, is_dir) = |
| 829 | store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap(); |
| 830 | assert!(is_dir); |
| 831 | let stored = dir.path().join("doc.tar.gz"); |
| 832 | assert_eq!(size, stored.metadata().unwrap().len() as i64); |
| 833 | // Round-trips through gzip back to the original tar bytes. |
| 834 | let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap()); |
| 835 | let mut bytes = Vec::new(); |
| 836 | gz.read_to_end(&mut bytes).unwrap(); |
| 837 | assert_eq!(bytes, tar); |
| 838 | } |
| 839 | |
| 840 | #[test] |
| 841 | fn extracts_a_browsable_directory_artifact() { |
| 842 | let dir = tempfile::tempdir().unwrap(); |
| 843 | let tar = tar_of(&[ |
| 844 | ("doc", None), |
| 845 | ("doc/index.html", Some("<html>")), |
| 846 | ("doc/sub", None), |
| 847 | ("doc/sub/page.html", Some("<p>")), |
| 848 | ]); |
| 849 | let (size, is_dir) = |
| 850 | store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap(); |
| 851 | assert!(is_dir); |
| 852 | assert_eq!(size, 6 + 3); |
| 853 | let root = dir.path().join("doc"); |
| 854 | assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>"); |
| 855 | assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>"); |
| 856 | } |
| 857 | |
| 858 | #[test] |
| 859 | fn parses_meta_tar_with_trimmed_capped_values() { |
| 860 | let tar = tar_of(&[ |
| 861 | ("anvil-meta", None), |
| 862 | ("anvil-meta/bin", None), |
| 863 | ("anvil-meta/bin/version", Some("anvild 0.0.0\n")), |
| 864 | ("anvil-meta/bin/empty", Some(" \n")), |
| 865 | ("anvil-meta/doc", None), |
| 866 | ("anvil-meta/doc/pages", Some("42")), |
| 867 | ]); |
| 868 | let metas = parse_meta_tar(&tar); |
| 869 | assert_eq!(metas["bin"]["version"], "anvild 0.0.0"); |
| 870 | assert_eq!(metas["doc"]["pages"], "42"); |
| 871 | assert!(!metas["bin"].contains_key("empty"), "blank values dropped"); |
| 872 | } |
| 873 | } |