anvilsign in

collin/anvil

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
16use std::{
17 collections::BTreeMap,
18 io::Read,
19 path::Path,
20};
21
22use anvil_core::{
23 App,
24 ci::{
25 self,
26 Pipeline,
27 },
28 config::CiConfig,
29 repos,
30 storage,
31 users,
32};
33use anvil_git::browse::{
34 self,
35 TreeFile,
36};
37use anvil_job::{
38 ArtifactSpec,
39 JobResult,
40 JobSpec,
41 META_DIR,
42 Sandbox,
43 Stored,
44 mb_cap,
45};
46use 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.
54fn 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.
115fn 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.
150pub 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.
187async 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.
212pub 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.
275fn 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.
305async 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.
323async 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.
421pub 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.
432pub 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.
460pub 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.
473fn 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.
482pub 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.
560async 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.
567async 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.
592async 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.
657fn 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).
663async 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.
694fn 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`
707fn 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.
781fn 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.
798fn single_quote(s: &str) -> String {
799 format!("'{}'", s.replace('\'', "'\\''"))
800}
801
802#[cfg(test)]
803mod 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}