anvilsign in

collin/anvil

main / crates / anvil-ci / src / lib.rs
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.toml`, 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.toml`, 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/// Read and parse the pipeline a run's commit defines.
301///
302/// The three callers below each need the whole `Pipeline` for a different
303/// field, so the read lives here rather than three times over.
304///
305/// A commit that has only the pre-TOML `.anvil/ci.yml` gets told to rename it.
306/// Without that the run would fail as "pipeline missing", which is true and
307/// useless — the file is right there.
308fn load_pipeline(repo_path: &Path, commit: &str) -> Result<ci::Pipeline, String> {
309 let src = browse::read_blob(repo_path, commit, ci::PIPELINE_PATH).map_err(|e| e.to_string())?;
310 let Some(src) = src else {
311 let legacy = browse::read_blob(repo_path, commit, ci::LEGACY_PIPELINE_PATH);
312 return Err(if matches!(legacy, Ok(Some(_))) {
313 format!(
314 "{} is no longer read — pipelines are TOML now. Rename it to {} and convert it \
315 (see docs/ci-artifacts.md for the shape).",
316 ci::LEGACY_PIPELINE_PATH,
317 ci::PIPELINE_PATH,
318 )
319 } else {
320 format!("{} missing at {commit}", ci::PIPELINE_PATH)
321 });
322 };
323 ci::parse_pipeline(&String::from_utf8_lossy(&src)).map_err(|e| e.to_string())
324}
325
326/// The platform a queued run needs, for the routing decision — `None` when it
327/// names none and can run anywhere.
328///
329/// Reads the pipeline and nothing else: no vault, no lease, no `mark_running`.
330/// [`prepare`] does that for the run that is actually taken.
331async fn target_platform(app: &App, run_id: i64) -> Result<Option<String>, String> {
332 let (_, _, repo_path, run) = resolve(app, run_id).await?;
333 let pipeline = load_pipeline(&repo_path, &run.commit)?;
334 Ok(app
335 .config
336 .ci
337 .resolve_platform(&pipeline.platform)
338 .map(str::to_string))
339}
340
341/// Everything anvil does before a job leaves the building: resolve the repo,
342/// parse the pipeline, open the vault, build the spec, take the lease.
343///
344/// `Ok(None)` means the run was finished here and should not be dispatched.
345async fn prepare(
346 app: &App,
347 run_id: i64,
348 runner: &str,
349 runner_platform: &str,
350) -> Result<Option<JobSpec>, String> {
351 let run = ci::get(&app.db, run_id)
352 .await
353 .map_err(|e| e.to_string())?
354 .ok_or("run not found")?;
355 let repo = repos::find_by_id(&app.db, run.repo_id)
356 .await
357 .map_err(|e| e.to_string())?
358 .ok_or("repository not found")?;
359 let owner = users::find_by_id(&app.db, repo.owner_id)
360 .await
361 .map_err(|e| e.to_string())?
362 .ok_or("owner not found")?;
363 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
364
365 let pipeline = load_pipeline(&repo_path, &run.commit)?;
366
367 let short = &run.commit[..run.commit.len().min(12)];
368 let mut log = format!(
369 "anvil ci · {}/{} · {} @ {short}\nrunner: {runner}\nimage: {}\n",
370 owner.username,
371 repo.name,
372 run.ref_name,
373 app.config.ci.resolve_image(&pipeline.image)
374 );
375 // Say so when the job is about to be emulated. Without this line a slow
376 // amd64-on-arm64 build looks like a slow build.
377 if let Some(platform) = app.config.ci.resolve_platform(&pipeline.platform) {
378 log.push_str(&format!("platform: {platform}"));
379 if !runner_platform.is_empty() && runner_platform != platform {
380 log.push_str(&format!(" (emulated on {runner_platform})"));
381 }
382 log.push('\n');
383 }
384
385 // Secrets the pipeline asked for, from the in-memory vault. anvil holds no
386 // key that opens the stored envelopes, so an unlock must have happened
387 // (`anvild secret unlock`) or the run cannot proceed.
388 let env = match app.vault.take(run.repo_id, &pipeline.secrets) {
389 Ok(values) => values,
390 Err(missing) => {
391 log.push_str(&format!(
392 "\n[secrets unavailable: {}]\n\
393 This repository is sealed or was unlocked without them. Run:\n \
394 anvild secret unlock {}/{}\n",
395 missing.join(", "),
396 owner.username,
397 repo.name,
398 ));
399 ci::append_log(&app.db, run_id, &log).await.ok();
400 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
401 tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name);
402 return Ok(None);
403 }
404 };
405 if !env.is_empty() {
406 log.push_str(&format!(
407 "secrets: {}\n",
408 env.iter()
409 .map(|(name, _)| name.as_str())
410 .collect::<Vec<_>>()
411 .join(", ")
412 ));
413 }
414
415 let job = match build_job(run_id, &pipeline, &app.config.ci, &env) {
416 Ok(job) => job,
417 Err(e) => {
418 log.push_str(&format!("\n[runner error] {e}\n"));
419 ci::append_log(&app.db, run_id, &log).await.ok();
420 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
421 return Ok(None);
422 }
423 };
424
425 // The header is written now rather than kept in memory until the result
426 // lands, so a run that is visibly `running` has something to show.
427 ci::append_log(&app.db, run_id, &log).await.ok();
428 ci::mark_running(&app.db, run_id).await.ok();
429 app.jobs.claim(run_id, runner, env);
430 tracing::info!("ci: run {run_id} claimed by {runner}");
431 Ok(Some(job))
432}
433
434/// The checkout a claimed job runs against, as an uncompressed tar.
435///
436/// Rebuilt from the commit on demand rather than stashed at claim time: it is
437/// a pure function of the commit, so a runner retrying the fetch is free, and
438/// anvil holds no per-job buffer on a host where memory is the scarce thing.
439pub async fn checkout_tar(app: &App, run_id: i64) -> Result<Vec<u8>, String> {
440 let (_, _, repo_path, run) = resolve(app, run_id).await?;
441 let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?;
442 Ok(build_tar(&files))
443}
444
445/// The declared spec for one artifact of a run.
446///
447/// Re-derived from the pipeline rather than trusted from the upload request:
448/// `browse` decides whether a tarball is extracted into a servable directory
449/// tree, which is not a choice a runner should get to make.
450pub async fn artifact_spec(
451 app: &App,
452 run_id: i64,
453 name: &str,
454) -> Result<Option<ArtifactSpec>, String> {
455 let (_, _, repo_path, run) = resolve(app, run_id).await?;
456 let pipeline = load_pipeline(&repo_path, &run.commit)?;
457 Ok(pipeline
458 .artifacts
459 .iter()
460 .find(|a| a.name == name)
461 .map(|a| ArtifactSpec {
462 name: a.name.clone(),
463 path: a.path.clone(),
464 browse: a.browse,
465 has_meta: !a.meta.is_empty(),
466 }))
467}
468
469/// Store one uploaded artifact tar into the run's scratch directory.
470///
471/// Scratch sits next to the artifacts' final home (same filesystem, so
472/// [`finish_run`]'s swap is a rename) and is only moved into place once the
473/// run reports in.
474pub fn store_upload(
475 app: &App,
476 run: &anvil_core::CiRun,
477 spec: &ArtifactSpec,
478 tar: &[u8],
479) -> Result<Stored, String> {
480 let scratch = scratch_dir(app, run);
481 std::fs::create_dir_all(&scratch).map_err(|e| format!("creating scratch dir: {e}"))?;
482 let (size, is_dir) = store_artifact(spec, tar, &scratch)?;
483 Ok(Stored { size, is_dir })
484}
485
486/// Where a run's artifacts accumulate before the swap.
487fn scratch_dir(app: &App, run: &anvil_core::CiRun) -> std::path::PathBuf {
488 app.config
489 .artifacts_dir()
490 .join(run.repo_id.to_string())
491 .join(format!(".collecting-{}", run.id))
492}
493
494/// Record a finished job: swap its artifacts into place, mask and store the
495/// log, set the status, and fire the deploy webhook if it qualifies.
496pub async fn finish_run(app: &App, run_id: i64, result: JobResult) -> Result<(), String> {
497 let (owner, repo, repo_path, run) = resolve(app, run_id).await?;
498 let scratch = scratch_dir(app, &run);
499 let mut log = result.log;
500
501 let (status, collected) = match result.runner_error {
502 Some(e) => {
503 log.push_str(&format!("\n[runner error] {e}\n"));
504 (ci::status::ERROR, Vec::new())
505 }
506 None if result.exit_code == 0 => (ci::status::SUCCESS, result.artifacts),
507 None => {
508 log.push_str(&format!("\n[exited with status {}]\n", result.exit_code));
509 (ci::status::FAILURE, result.artifacts)
510 }
511 };
512
513 // Swap the collected set into place, replacing any earlier run's
514 // artifacts for this commit, then record the rows.
515 if !collected.is_empty() {
516 let artifacts_root = app.config.artifacts_dir();
517 let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit);
518 let swap = async {
519 ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit)
520 .await
521 .map_err(|e| e.to_string())?;
522 if final_dir.exists() {
523 std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?;
524 }
525 std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?;
526 for c in &collected {
527 ci::add_artifact(
528 &app.db,
529 run_id,
530 run.repo_id,
531 &run.commit,
532 &c.name,
533 c.size,
534 c.is_dir,
535 c.browse,
536 &c.meta,
537 )
538 .await
539 .map_err(|e| e.to_string())?;
540 }
541 Ok::<_, String>(())
542 };
543 match swap.await {
544 Ok(()) => {
545 log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len()));
546 gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await;
547 }
548 Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")),
549 }
550 }
551 let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away
552
553 // A step that echoes its environment (`set -x`, `curl -v`, a failing
554 // command that prints its arguments) would otherwise publish the value on
555 // a page anyone with read access can see. The values come from the lease
556 // rather than the vault, which may have re-sealed while the job ran.
557 mask_secrets(&mut log, &app.jobs.secrets_for(run_id));
558 ci::append_log(&app.db, run_id, &log).await.ok();
559 ci::finish(&app.db, run_id, status).await.ok();
560 app.jobs.release(run_id);
561 tracing::info!("ci: run {run_id} {status}");
562
563 // Continuous deployment: on a green run of the configured deploy repo's
564 // deploy branch, fire the redeploy webhook. Scoped to one repo by config —
565 // no other repository can trigger it, even with passing CI.
566 if status == ci::status::SUCCESS && app.config.ci.is_deploy_target(&owner, &repo, &run.ref_name)
567 {
568 deploy(app, &owner, &repo, &run).await;
569 }
570 Ok(())
571}
572
573/// Finish a run that never got far enough to produce a result.
574async fn fail_run(app: &App, run_id: i64, note: &str) {
575 ci::append_log(&app.db, run_id, note).await.ok();
576 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
577 app.jobs.release(run_id);
578}
579
580/// (owner username, repo name, repo path, run) for a run id.
581async fn resolve(
582 app: &App,
583 run_id: i64,
584) -> Result<(String, String, std::path::PathBuf, anvil_core::CiRun), String> {
585 let run = ci::get(&app.db, run_id)
586 .await
587 .map_err(|e| e.to_string())?
588 .ok_or("run not found")?;
589 let repo = repos::find_by_id(&app.db, run.repo_id)
590 .await
591 .map_err(|e| e.to_string())?
592 .ok_or("repository not found")?;
593 let owner = users::find_by_id(&app.db, repo.owner_id)
594 .await
595 .map_err(|e| e.to_string())?
596 .ok_or("owner not found")?;
597 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
598 Ok((owner.username, repo.name, repo_path, run))
599}
600
601/// Enforce `[ci] artifact_quota_mb` for one repository: while over budget,
602/// delete the oldest commit's artifacts (rows + directory). Branch-tip
603/// commits and the just-stored commit are pinned. Deterministic — runs after
604/// every artifact-producing run, no background sweeper. Best-effort: failures
605/// are logged, never failing the run.
606async fn gc_artifacts(
607 app: &App,
608 repo_id: i64,
609 repo_path: &Path,
610 keep_commit: &str,
611 log: &mut String,
612) {
613 let quota = match app.config.ci.artifact_quota_mb {
614 0 => return,
615 mb => mb_cap(mb),
616 };
617 let rows = match ci::artifacts_for_repo(&app.db, repo_id).await {
618 Ok(rows) => rows,
619 Err(e) => {
620 log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n"));
621 return;
622 }
623 };
624
625 // Per-commit totals and ages.
626 let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes)
627 let mut total: u64 = 0;
628 for row in &rows {
629 let entry = commits
630 .entry(row.commit.clone())
631 .or_insert((row.created_at, 0));
632 entry.0 = entry.0.min(row.created_at);
633 entry.1 += row.size;
634 total = total.saturating_add(row.size.max(0) as u64);
635 }
636 if total <= quota {
637 return;
638 }
639
640 let pinned = browse::branch_tips(repo_path).unwrap_or_default();
641 let mut victims: Vec<(i64, String, i64)> = commits
642 .into_iter()
643 .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit))
644 .map(|(commit, (oldest, bytes))| (oldest, commit, bytes))
645 .collect();
646 victims.sort();
647
648 for (_, commit, bytes) in victims {
649 if total <= quota {
650 break;
651 }
652 let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit);
653 if let Err(e) = std::fs::remove_dir_all(&dir) {
654 log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n"));
655 continue; // keep the rows; retried next run
656 }
657 if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await {
658 log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n"));
659 continue;
660 }
661 total = total.saturating_sub(bytes.max(0) as u64);
662 log.push_str(&format!(
663 "\n[artifact gc: dropped {} for old commit {}]\n",
664 fmt_mb(bytes),
665 &commit[..commit.len().min(12)]
666 ));
667 }
668}
669
670/// Bytes as a short MiB string for log lines.
671fn fmt_mb(bytes: i64) -> String {
672 format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0))
673}
674
675/// POST the configured deploy webhook. Best-effort: logs success/failure but
676/// never fails the run (CI already passed).
677async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) {
678 let cfg = &app.config.ci;
679 let body = serde_json::json!({
680 "repo": format!("{owner}/{name}"),
681 "ref": run.ref_name,
682 "commit": run.commit,
683 "run_id": run.id,
684 });
685 let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body);
686 if !cfg.deploy_secret.is_empty() {
687 req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret);
688 }
689 match req.send().await {
690 Ok(resp) if resp.status().is_success() => {
691 tracing::info!(
692 "ci: deploy webhook for {owner}/{name} accepted ({})",
693 resp.status()
694 )
695 }
696 Ok(resp) => tracing::error!(
697 "ci: deploy webhook for {owner}/{name} returned {}",
698 resp.status()
699 ),
700 Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"),
701 }
702}
703
704/// Replace every secret value in `log` with `***`.
705///
706/// Only values worth hiding: very short ones (a one-character secret) would
707/// mask half the log for no benefit, and are not credentials in practice.
708fn mask_secrets(log: &mut String, env: &[(String, String)]) {
709 for (_, value) in env {
710 if value.len() >= 4 && log.contains(value.as_str()) {
711 *log = log.replace(value.as_str(), "***");
712 }
713 }
714}
715
716/// Write one downloaded artifact tar into `scratch`, returning (size, is_dir).
717///
718/// - a file is stored as-is at `scratch/<name>`
719/// - a directory with `browse` is extracted under `scratch/<name>/`
720/// - any other directory is stored compressed at `scratch/<name>.tar.gz`
721fn store_artifact(
722 spec: &ArtifactSpec,
723 tar_bytes: &[u8],
724 scratch: &Path,
725) -> Result<(i64, bool), String> {
726 // The docker archive endpoint roots entries at the requested item's
727 // basename; its first entry tells file from directory.
728 let mut archive = tar::Archive::new(tar_bytes);
729 let mut entries = archive.entries().map_err(|e| e.to_string())?;
730 let first = entries
731 .next()
732 .ok_or("empty archive")?
733 .map_err(|e| e.to_string())?;
734 let is_dir = first.header().entry_type().is_dir();
735
736 if !is_dir {
737 let mut entry = first;
738 let mut content = Vec::new();
739 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
740 std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?;
741 return Ok((content.len() as i64, false));
742 }
743
744 if !spec.browse {
745 let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name)))
746 .map_err(|e| e.to_string())?;
747 let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default());
748 std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?;
749 let file = enc.finish().map_err(|e| e.to_string())?;
750 let size = file.metadata().map_err(|e| e.to_string())?.len();
751 return Ok((size as i64, true));
752 }
753
754 // Browsable: extract regular files under scratch/<name>/, stripping the
755 // basename prefix. Entry paths come from docker's tar of a real
756 // filesystem, but stay defensive: relative components only, no links.
757 let root = scratch.join(&spec.name);
758 let mut total = 0i64;
759 let mut archive = tar::Archive::new(tar_bytes);
760 for entry in archive.entries().map_err(|e| e.to_string())?.flatten() {
761 if !entry.header().entry_type().is_file() {
762 continue;
763 }
764 let Ok(path) = entry.path() else { continue };
765 let mut rel = std::path::PathBuf::new();
766 let mut ok = true;
767 for c in path.components().skip(1) {
768 match c {
769 std::path::Component::Normal(p) => rel.push(p),
770 _ => {
771 ok = false;
772 break;
773 }
774 }
775 }
776 if !ok || rel.as_os_str().is_empty() {
777 continue;
778 }
779 let dest = root.join(&rel);
780 if let Some(parent) = dest.parent() {
781 std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
782 }
783 let mut entry = entry;
784 let mut content = Vec::new();
785 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
786 std::fs::write(&dest, &content).map_err(|e| e.to_string())?;
787 total += content.len() as i64;
788 }
789 let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists
790 Ok((total, true))
791}
792
793/// Build an uncompressed tar of the checkout, rooted at `workspace/` so it
794/// extracts to `/workspace` when uploaded to the container root.
795fn build_tar(files: &[TreeFile]) -> Vec<u8> {
796 let mut builder = tar::Builder::new(Vec::new());
797 for f in files {
798 let mut header = tar::Header::new_gnu();
799 header.set_size(f.content.len() as u64);
800 header.set_mode(if f.executable { 0o755 } else { 0o644 });
801 // append_data sets the path and checksum.
802 let _ = builder.append_data(
803 &mut header,
804 format!("workspace/{}", f.path),
805 f.content.as_slice(),
806 );
807 }
808 builder.into_inner().unwrap_or_default()
809}
810
811/// Single-quote a string for safe interpolation into a shell command.
812fn single_quote(s: &str) -> String {
813 format!("'{}'", s.replace('\'', "'\\''"))
814}
815
816#[cfg(test)]
817mod tests {
818 use super::*;
819
820 /// Build a tar the way docker's archive endpoint does: entries rooted at
821 /// the requested item's basename.
822 fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> {
823 let mut b = tar::Builder::new(Vec::new());
824 for (path, content) in entries {
825 let mut h = tar::Header::new_gnu();
826 match content {
827 Some(c) => {
828 h.set_size(c.len() as u64);
829 h.set_mode(0o644);
830 h.set_entry_type(tar::EntryType::Regular);
831 b.append_data(&mut h, path, c.as_bytes()).unwrap();
832 }
833 None => {
834 h.set_size(0);
835 h.set_mode(0o755);
836 h.set_entry_type(tar::EntryType::Directory);
837 b.append_data(&mut h, path, std::io::empty()).unwrap();
838 }
839 }
840 }
841 b.into_inner().unwrap()
842 }
843
844 /// Queued runs as (id, wanted platform), oldest first.
845 fn queue(items: &[(i64, Option<&str>)]) -> Vec<(i64, Option<String>)> {
846 items
847 .iter()
848 .map(|(id, p)| (*id, p.map(str::to_string)))
849 .collect()
850 }
851
852 /// A runner is offered its own platform's jobs and the unconstrained ones,
853 /// in queue order, ahead of anything it would have to emulate.
854 #[test]
855 fn claim_order_prefers_native_work() {
856 let queued = queue(&[
857 (1, Some("linux/amd64")),
858 (2, None),
859 (3, Some("linux/arm64")),
860 (4, Some("linux/amd64")),
861 ]);
862 let live = ["linux/amd64".to_string(), "linux/arm64".to_string()];
863
864 // Both architectures are present, so neither runner touches the
865 // other's jobs — not even when its own queue is empty.
866 assert_eq!(claim_order(&queued, "linux/arm64", &live), vec![2, 3]);
867 assert_eq!(claim_order(&queued, "linux/amd64", &live), vec![1, 2, 4]);
868 }
869
870 /// With no runner of the named architecture present, the job is offered to
871 /// whoever asks rather than sitting queued forever — after that runner's
872 /// own work. This is the single-runner case: one arm64 Mac, pipelines that
873 /// ask for amd64 because that is what they ship.
874 #[test]
875 fn claim_order_falls_back_to_emulation_when_nobody_is_native() {
876 let queued = queue(&[(1, Some("linux/amd64")), (2, None)]);
877 let live = ["linux/arm64".to_string()];
878 assert_eq!(claim_order(&queued, "linux/arm64", &live), vec![2, 1]);
879 }
880
881 /// An unknown platform (`linux/riscv64`, a typo) behaves like any other
882 /// platform nobody is native to: it runs somewhere and fails visibly there,
883 /// rather than disappearing from the queue.
884 #[test]
885 fn claim_order_never_strands_a_run() {
886 let queued = queue(&[(1, Some("linux/riscv64"))]);
887 assert_eq!(
888 claim_order(&queued, "linux/arm64", &["linux/arm64".to_string()]),
889 vec![1]
890 );
891 }
892
893 fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec {
894 ArtifactSpec {
895 name: name.into(),
896 path: path.into(),
897 browse,
898 has_meta: false,
899 }
900 }
901
902 #[test]
903 fn stores_a_file_artifact_as_is() {
904 let dir = tempfile::tempdir().unwrap();
905 let tar = tar_of(&[("anvild", Some("ELF..."))]);
906 let (size, is_dir) = store_artifact(
907 &spec("bin", "target/release/anvild", false),
908 &tar,
909 dir.path(),
910 )
911 .unwrap();
912 assert!(!is_dir);
913 assert_eq!(size, 6);
914 assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF...");
915 }
916
917 #[test]
918 fn stores_a_directory_artifact_as_tar_gz() {
919 let dir = tempfile::tempdir().unwrap();
920 let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]);
921 let (size, is_dir) =
922 store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap();
923 assert!(is_dir);
924 let stored = dir.path().join("doc.tar.gz");
925 assert_eq!(size, stored.metadata().unwrap().len() as i64);
926 // Round-trips through gzip back to the original tar bytes.
927 let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap());
928 let mut bytes = Vec::new();
929 gz.read_to_end(&mut bytes).unwrap();
930 assert_eq!(bytes, tar);
931 }
932
933 #[test]
934 fn extracts_a_browsable_directory_artifact() {
935 let dir = tempfile::tempdir().unwrap();
936 let tar = tar_of(&[
937 ("doc", None),
938 ("doc/index.html", Some("<html>")),
939 ("doc/sub", None),
940 ("doc/sub/page.html", Some("<p>")),
941 ]);
942 let (size, is_dir) =
943 store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap();
944 assert!(is_dir);
945 assert_eq!(size, 6 + 3);
946 let root = dir.path().join("doc");
947 assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>");
948 assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>");
949 }
950}