anvilsign in

collin/anvil

1//! Continuous integration: the pipeline definition (`.anvil/ci.toml`) and the
2//! persistence/lifecycle of CI runs. Execution (Docker) lives in `anvil-ci`.
3
4use serde::Deserialize;
5
6use crate::{
7 error::{
8 Error,
9 Result,
10 },
11 models::{
12 CiArtifact,
13 CiRun,
14 },
15};
16
17/// Run status values stored in [`CiRun::status`].
18pub mod status {
19 pub const QUEUED: &str = "queued";
20 pub const RUNNING: &str = "running";
21 pub const SUCCESS: &str = "success";
22 pub const FAILURE: &str = "failure";
23 pub const ERROR: &str = "error";
24}
25
26/// Path of the pipeline definition within a repository.
27pub const PIPELINE_PATH: &str = ".anvil/ci.yml";
28
29/// A parsed pipeline: a base image, ordered straight-line steps, and the
30/// artifacts to collect afterwards.
31#[derive(Clone, Debug, Deserialize)]
32pub struct Pipeline {
33 /// Docker image the steps run in, e.g. `rust:1.95-bookworm`. Optional:
34 /// omitting it selects `ci.default_image`, the shared anvil runner that
35 /// agent sessions also use. Resolve it with
36 /// [`CiConfig::resolve_image`](crate::config::CiConfig::resolve_image)
37 /// rather than reading this field directly — it is empty when omitted.
38 #[serde(default)]
39 pub image: String,
40 /// Docker platform to run on, `os/arch[/variant]` — e.g. `linux/amd64` on
41 /// a repository that ships amd64 but has arm64 runners. Empty falls back to
42 /// `[ci] platform`, and if that is empty too, to whatever the claiming
43 /// runner's daemon is native to. Resolve it with
44 /// [`CiConfig::resolve_platform`](crate::config::CiConfig::resolve_platform)
45 /// rather than reading this field.
46 ///
47 /// It is both an execution setting and a scheduling one: anvil hands the
48 /// job to a runner that is natively this platform when it has one, and
49 /// only falls back to an emulating runner when it does not.
50 #[serde(default)]
51 pub platform: String,
52 #[serde(default)]
53 pub steps: Vec<Step>,
54 #[serde(default)]
55 pub artifacts: Vec<ArtifactSpec>,
56 /// Names of repository secrets to expose as environment variables (see
57 /// `docs/secrets.md`). The run fails before starting a container unless
58 /// every one of them is available, which requires the repository to be
59 /// unlocked — anvil cannot decrypt them by itself.
60 #[serde(default)]
61 pub secrets: Vec<String>,
62}
63
64/// A declared artifact: a path in the workspace to collect after the steps
65/// run, plus optional metadata extractors (see `docs/ci-artifacts.md`).
66#[derive(Clone, Debug, Deserialize)]
67pub struct ArtifactSpec {
68 /// Display/URL name; unique within the pipeline, `[A-Za-z0-9._-]+`.
69 pub name: String,
70 /// Path relative to the workspace root. A file downloads as-is; a
71 /// directory downloads as a tarball — unless `browse` is set.
72 pub path: String,
73 /// Serve this (directory) artifact as a browsable static site instead of
74 /// a download, e.g. rustdoc output.
75 #[serde(default)]
76 pub browse: bool,
77 /// Metadata extractors: key → shell command, run *inside the job
78 /// container* after the steps. Each command's stdout (trimmed, capped)
79 /// becomes the value shown next to the artifact.
80 #[serde(default)]
81 pub meta: std::collections::BTreeMap<String, String>,
82}
83
84/// One pipeline step: a shell command, with an optional display name.
85#[derive(Clone, Debug, Deserialize)]
86pub struct Step {
87 #[serde(default)]
88 pub name: String,
89 pub run: String,
90}
91
92impl Step {
93 /// Display label: the explicit name, or the command if unnamed.
94 pub fn label(&self) -> &str {
95 if self.name.is_empty() {
96 &self.run
97 } else {
98 &self.name
99 }
100 }
101}
102
103/// Whether `platform` is a well-formed Docker platform: `os/arch`, optionally
104/// with a variant (`linux/arm/v7`), all lowercase `[a-z0-9._-]`.
105///
106/// Shape only, deliberately: the set of valid architectures is the daemon's to
107/// know, and a forge that hardcoded one would need a release to gain `riscv64`.
108/// What this does catch is the mistake worth catching — a bare `amd64` with no
109/// `os/`, which Docker would silently read as an operating system.
110pub fn valid_platform(platform: &str) -> bool {
111 let parts: Vec<&str> = platform.split('/').collect();
112 (2..=3).contains(&parts.len())
113 && parts.iter().all(|part| {
114 !part.is_empty()
115 && part
116 .chars()
117 .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || ".-_".contains(c))
118 })
119}
120
121/// Parse a `.anvil/ci.yml` pipeline definition.
122pub fn parse_pipeline(yaml: &str) -> Result<Pipeline> {
123 let mut pipeline: Pipeline = serde_yaml::from_str(yaml)
124 .map_err(|e| Error::Invalid(format!("invalid {PIPELINE_PATH}: {e}")))?;
125 if pipeline.image.trim().is_empty() {
126 return Err(Error::Invalid(format!(
127 "{PIPELINE_PATH}: `image` is required"
128 )));
129 }
130 pipeline.platform = pipeline.platform.trim().to_string();
131 if !pipeline.platform.is_empty() && !valid_platform(&pipeline.platform) {
132 return Err(Error::Invalid(format!(
133 "{PIPELINE_PATH}: platform `{}` must be os/arch, e.g. linux/amd64",
134 pipeline.platform
135 )));
136 }
137 let mut seen = std::collections::BTreeSet::new();
138 for a in &mut pipeline.artifacts {
139 let invalid =
140 |what: &str| Error::Invalid(format!("{PIPELINE_PATH}: artifact `{}`: {what}", a.name));
141 if a.name.is_empty()
142 || !a
143 .name
144 .chars()
145 .all(|c| c.is_ascii_alphanumeric() || ".-_".contains(c))
146 {
147 return Err(invalid("name must be non-empty [A-Za-z0-9._-]+"));
148 }
149 if !seen.insert(a.name.clone()) {
150 return Err(invalid("duplicate name"));
151 }
152 // Normalize away a trailing slash so a directory path and the same
153 // path without the slash behave identically downstream.
154 a.path = a.path.trim_end_matches('/').to_string();
155 if a.path.is_empty()
156 || a.path.starts_with('/')
157 || a.path.split('/').any(|seg| seg == ".." || seg.is_empty())
158 {
159 return Err(invalid(
160 "path must be relative to the workspace, without `..`",
161 ));
162 }
163 // Meta keys become file names in the extractor handoff, so they get
164 // the same charset as artifact names.
165 for key in a.meta.keys() {
166 if key.is_empty()
167 || !key
168 .chars()
169 .all(|c| c.is_ascii_alphanumeric() || ".-_".contains(c))
170 {
171 return Err(invalid(&format!(
172 "meta key `{key}` must be [A-Za-z0-9._-]+"
173 )));
174 }
175 }
176 }
177 for name in &pipeline.secrets {
178 if !crate::secrets::valid_name(name) {
179 return Err(Error::Invalid(format!(
180 "{PIPELINE_PATH}: secret `{name}` must be A–Z, 0–9 and _, not starting with a digit"
181 )));
182 }
183 }
184 Ok(pipeline)
185}
186
187/// Create a queued CI run for a pushed commit.
188pub async fn enqueue(db: &toasty::Db, repo_id: i64, commit: &str, ref_name: &str) -> Result<CiRun> {
189 let mut conn = db.clone();
190 let run = toasty::create!(CiRun {
191 repo_id: repo_id,
192 commit: commit,
193 ref_name: ref_name,
194 status: status::QUEUED,
195 log: "",
196 created_at: crate::now(),
197 started_at: 0,
198 finished_at: 0,
199 })
200 .exec(&mut conn)
201 .await?;
202 Ok(run)
203}
204
205/// Fetch a run by id.
206pub async fn get(db: &toasty::Db, id: i64) -> Result<Option<CiRun>> {
207 let mut conn = db.clone();
208 Ok(CiRun::filter(CiRun::fields().id().eq(id))
209 .first()
210 .exec(&mut conn)
211 .await?)
212}
213
214/// List a repository's runs, newest first, up to `limit`.
215pub async fn list_by_repo(db: &toasty::Db, repo_id: i64, limit: usize) -> Result<Vec<CiRun>> {
216 let mut conn = db.clone();
217 let runs = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
218 .order_by(CiRun::fields().id().desc())
219 .limit(limit)
220 .exec(&mut conn)
221 .await?;
222 Ok(runs)
223}
224
225/// The most recent run for a specific commit (for status badges).
226pub async fn latest_for_commit(
227 db: &toasty::Db,
228 repo_id: i64,
229 commit: &str,
230) -> Result<Option<CiRun>> {
231 let mut conn = db.clone();
232 let run = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
233 .filter(CiRun::fields().commit().eq(commit))
234 .order_by(CiRun::fields().id().desc())
235 .first()
236 .exec(&mut conn)
237 .await?;
238 Ok(run)
239}
240
241/// Mark a run as started (running).
242pub async fn mark_running(db: &toasty::Db, id: i64) -> Result<()> {
243 let Some(mut run) = get(db, id).await? else {
244 return Ok(());
245 };
246 let mut conn = db.clone();
247 run.update()
248 .status(status::RUNNING)
249 .started_at(crate::now())
250 .exec(&mut conn)
251 .await?;
252 Ok(())
253}
254
255/// Append a chunk to a run's log.
256pub async fn append_log(db: &toasty::Db, id: i64, chunk: &str) -> Result<()> {
257 let Some(mut run) = get(db, id).await? else {
258 return Ok(());
259 };
260 let combined = format!("{}{chunk}", run.log);
261 let mut conn = db.clone();
262 run.update().log(&combined).exec(&mut conn).await?;
263 Ok(())
264}
265
266/// Finish a run with a terminal status (`success`/`failure`/`error`).
267pub async fn finish(db: &toasty::Db, id: i64, status: &str) -> Result<()> {
268 let Some(mut run) = get(db, id).await? else {
269 return Ok(());
270 };
271 let mut conn = db.clone();
272 run.update()
273 .status(status)
274 .finished_at(crate::now())
275 .exec(&mut conn)
276 .await?;
277 Ok(())
278}
279
280/// Re-queue runs left mid-flight by a crash/restart (status `running`).
281/// Returns the ids that were requeued so the runner can pick them up.
282pub async fn requeue_interrupted(db: &toasty::Db) -> Result<Vec<i64>> {
283 let mut conn = db.clone();
284 let interrupted = CiRun::filter(CiRun::fields().status().eq(status::RUNNING))
285 .exec(&mut conn)
286 .await?;
287 let mut ids = Vec::new();
288 for mut run in interrupted {
289 let mut conn = db.clone();
290 run.update().status(status::QUEUED).exec(&mut conn).await?;
291 ids.push(run.id);
292 }
293 Ok(ids)
294}
295
296/// Put one run back on the queue.
297///
298/// Used when a runner's lease expires: the job may still be executing on an
299/// unreachable machine, but from anvil's side it is no longer accounted for,
300/// and leaving it `running` forever is worse than the chance of a second
301/// attempt.
302pub async fn requeue(db: &toasty::Db, id: i64) -> Result<()> {
303 let mut conn = db.clone();
304 let Some(mut run) = get(db, id).await? else {
305 return Ok(());
306 };
307 run.update().status(status::QUEUED).exec(&mut conn).await?;
308 Ok(())
309}
310
311/// List all queued run ids (oldest first) — used on startup to drain the queue.
312pub async fn queued_ids(db: &toasty::Db) -> Result<Vec<i64>> {
313 let mut conn = db.clone();
314 let runs = CiRun::filter(CiRun::fields().status().eq(status::QUEUED))
315 .order_by(CiRun::fields().id().asc())
316 .exec(&mut conn)
317 .await?;
318 Ok(runs.into_iter().map(|r| r.id).collect())
319}
320
321/// Record a collected artifact. `meta` is a JSON object string (`{}` if none).
322#[allow(clippy::too_many_arguments)]
323pub async fn add_artifact(
324 db: &toasty::Db,
325 run_id: i64,
326 repo_id: i64,
327 commit: &str,
328 name: &str,
329 size: i64,
330 is_dir: bool,
331 browse: bool,
332 meta: &str,
333) -> Result<CiArtifact> {
334 let mut conn = db.clone();
335 let artifact = toasty::create!(CiArtifact {
336 run_id: run_id,
337 repo_id: repo_id,
338 commit: commit,
339 name: name,
340 size: size,
341 is_dir: is_dir,
342 browse: browse,
343 meta: meta,
344 created_at: crate::now(),
345 })
346 .exec(&mut conn)
347 .await?;
348 Ok(artifact)
349}
350
351/// A run's artifacts, in declaration (insertion) order.
352pub async fn artifacts_for_run(db: &toasty::Db, run_id: i64) -> Result<Vec<CiArtifact>> {
353 let mut conn = db.clone();
354 let artifacts = CiArtifact::filter(CiArtifact::fields().run_id().eq(run_id))
355 .order_by(CiArtifact::fields().id().asc())
356 .exec(&mut conn)
357 .await?;
358 Ok(artifacts)
359}
360
361/// The newest artifact named `name` for `commit` (across that commit's runs).
362pub async fn latest_artifact(
363 db: &toasty::Db,
364 repo_id: i64,
365 commit: &str,
366 name: &str,
367) -> Result<Option<CiArtifact>> {
368 let mut conn = db.clone();
369 let artifact = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
370 .filter(CiArtifact::fields().commit().eq(commit))
371 .filter(CiArtifact::fields().name().eq(name))
372 .order_by(CiArtifact::fields().id().desc())
373 .first()
374 .exec(&mut conn)
375 .await?;
376 Ok(artifact)
377}
378
379/// All artifact rows for a repository (for GC accounting). Small at our scale.
380pub async fn artifacts_for_repo(db: &toasty::Db, repo_id: i64) -> Result<Vec<CiArtifact>> {
381 let mut conn = db.clone();
382 let artifacts = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
383 .exec(&mut conn)
384 .await?;
385 Ok(artifacts)
386}
387
388/// Delete the artifact rows for one commit (the on-disk directory is the
389/// caller's to remove). Used when re-running a commit and by GC.
390pub async fn delete_artifacts_for_commit(
391 db: &toasty::Db,
392 repo_id: i64,
393 commit: &str,
394) -> Result<()> {
395 let mut conn = db.clone();
396 let rows = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
397 .filter(CiArtifact::fields().commit().eq(commit))
398 .exec(&mut conn)
399 .await?;
400 for row in rows {
401 let mut conn = db.clone();
402 row.delete().exec(&mut conn).await?;
403 }
404 Ok(())
405}
406
407/// Delete every run and artifact row belonging to a repository, returning the
408/// run ids that were removed so the caller can release their leases. On-disk
409/// artifact directories are the caller's to remove — see
410/// [`repos::delete`](crate::repos::delete), the only user.
411pub async fn delete_for_repo(db: &toasty::Db, repo_id: i64) -> Result<Vec<i64>> {
412 for artifact in artifacts_for_repo(db, repo_id).await? {
413 let mut conn = db.clone();
414 artifact.delete().exec(&mut conn).await?;
415 }
416 let mut conn = db.clone();
417 let runs = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
418 .exec(&mut conn)
419 .await?;
420 let mut ids = Vec::with_capacity(runs.len());
421 for run in runs {
422 ids.push(run.id);
423 let mut conn = db.clone();
424 run.delete().exec(&mut conn).await?;
425 }
426 Ok(ids)
427}
428
429#[cfg(test)]
430mod tests {
431 use super::*;
432
433 #[test]
434 fn parses_a_basic_pipeline() {
435 // YAML is indentation-sensitive, so the fixture is flush-left.
436 let p = parse_pipeline(
437 r#"image: rust:1.95-bookworm
438steps:
439 - name: test
440 run: cargo test --workspace
441 - run: cargo build --release
442"#,
443 )
444 .unwrap();
445 assert_eq!(p.image, "rust:1.95-bookworm");
446 assert_eq!(p.steps.len(), 2);
447 assert_eq!(p.steps[0].label(), "test");
448 // Unnamed step falls back to its command for the label.
449 assert_eq!(p.steps[1].label(), "cargo build --release");
450 }
451
452 #[test]
453 fn requires_an_image() {
454 assert!(parse_pipeline("steps: []\n").is_err());
455 }
456
457 /// `platform:` is optional, and when present has to be `os/arch` — a bare
458 /// `amd64` would otherwise reach Docker as an *operating system* named
459 /// amd64, which fails much later and much less clearly.
460 #[test]
461 fn parses_and_validates_the_platform() {
462 let p = parse_pipeline("image: rust:1.95\nplatform: linux/amd64\nsteps: []\n").unwrap();
463 assert_eq!(p.platform, "linux/amd64");
464 assert!(
465 parse_pipeline("image: rust:1.95\nsteps: []\n")
466 .unwrap()
467 .platform
468 .is_empty()
469 );
470
471 assert!(valid_platform("linux/arm64"));
472 assert!(valid_platform("linux/arm/v7"));
473 assert!(!valid_platform("amd64"));
474 assert!(!valid_platform("linux/"));
475 assert!(!valid_platform("linux/amd64/v1/extra"));
476 assert!(!valid_platform("Linux/AMD64"));
477 assert!(parse_pipeline("image: rust:1.95\nplatform: amd64\nsteps: []\n").is_err());
478 }
479
480 #[test]
481 fn parses_and_validates_artifacts() {
482 let p = parse_pipeline(
483 r#"image: rust:1.95
484artifacts:
485 - name: anvild
486 path: target/release/anvild
487 meta:
488 version: ./target/release/anvild --version
489 - name: doc
490 path: target/doc/
491 browse: true
492"#,
493 )
494 .unwrap();
495 assert_eq!(p.artifacts.len(), 2);
496 assert_eq!(p.artifacts[0].name, "anvild");
497 assert_eq!(
498 p.artifacts[0].meta["version"],
499 "./target/release/anvild --version"
500 );
501 assert!(!p.artifacts[0].browse);
502 assert_eq!(p.artifacts[1].path, "target/doc", "trailing slash trimmed");
503 assert!(p.artifacts[1].browse);
504
505 let must_fail = |yaml: &str, why: &str| {
506 assert!(
507 parse_pipeline(&format!("image: i\n{yaml}")).is_err(),
508 "{why}"
509 );
510 };
511 must_fail(
512 "artifacts: [{name: 'a b', path: x}]",
513 "space in name rejected",
514 );
515 must_fail(
516 "artifacts: [{name: 'a/b', path: x}]",
517 "slash in name rejected",
518 );
519 must_fail("artifacts: [{name: '', path: x}]", "empty name rejected");
520 must_fail(
521 "artifacts: [{name: a, path: x}, {name: a, path: y}]",
522 "duplicate name rejected",
523 );
524 must_fail(
525 "artifacts: [{name: a, path: /etc}]",
526 "absolute path rejected",
527 );
528 must_fail(
529 "artifacts: [{name: a, path: '../up'}]",
530 "traversal rejected",
531 );
532 must_fail(
533 "artifacts: [{name: a, path: 'x//y'}]",
534 "empty segment rejected",
535 );
536 must_fail("artifacts: [{name: a, path: ''}]", "empty path rejected");
537 }
538
539 // Exercises the run-lifecycle queries against a real SQLite database, to
540 // confirm the ORM-level `order_by`/`limit`/filter actually work (Toasty is
541 // pre-1.0, so we don't take that on faith).
542 #[tokio::test]
543 async fn ordering_and_limit_run_in_the_database() {
544 let dir = tempfile::tempdir().unwrap();
545 let db = crate::db::connect(dir.path().join("t.db")).await.unwrap();
546
547 for c in ["aaa", "bbb", "ccc"] {
548 enqueue(&db, 1, c, "main").await.unwrap();
549 }
550 enqueue(&db, 2, "zzz", "main").await.unwrap(); // a different repo
551
552 // Newest-first, limited to 2 — and scoped to repo 1.
553 let runs = list_by_repo(&db, 1, 2).await.unwrap();
554 assert_eq!(runs.len(), 2);
555 assert!(runs[0].id > runs[1].id, "newest first");
556 assert_eq!(runs[0].commit, "ccc");
557 assert!(runs.iter().all(|r| r.repo_id == 1));
558
559 // Commit filter pushed into the query (not an in-memory retain).
560 let latest = latest_for_commit(&db, 1, "bbb").await.unwrap().unwrap();
561 assert_eq!(latest.commit, "bbb");
562 assert_eq!(latest.repo_id, 1);
563 assert!(latest_for_commit(&db, 1, "nope").await.unwrap().is_none());
564
565 // All queued, ascending by id.
566 let queued = queued_ids(&db).await.unwrap();
567 assert_eq!(queued.len(), 4);
568 assert!(queued.windows(2).all(|w| w[0] < w[1]), "ascending");
569 }
570
571 #[tokio::test]
572 async fn artifact_rows_round_trip() {
573 let dir = tempfile::tempdir().unwrap();
574 let db = crate::db::connect(dir.path().join("t.db")).await.unwrap();
575
576 add_artifact(&db, 1, 7, "abc", "bin", 100, false, false, "{}")
577 .await
578 .unwrap();
579 add_artifact(&db, 1, 7, "abc", "doc", 5000, true, true, r#"{"v":"1"}"#)
580 .await
581 .unwrap();
582 // A newer run of the same commit supersedes for `latest_artifact`.
583 add_artifact(&db, 2, 7, "abc", "bin", 200, false, false, "{}")
584 .await
585 .unwrap();
586
587 let run1 = artifacts_for_run(&db, 1).await.unwrap();
588 assert_eq!(run1.len(), 2);
589 assert_eq!(run1[0].name, "bin");
590 assert!(run1[1].browse);
591
592 let latest = latest_artifact(&db, 7, "abc", "bin")
593 .await
594 .unwrap()
595 .unwrap();
596 assert_eq!(latest.run_id, 2);
597 assert_eq!(latest.size, 200);
598 assert!(
599 latest_artifact(&db, 8, "abc", "bin")
600 .await
601 .unwrap()
602 .is_none(),
603 "scoped to repo"
604 );
605
606 delete_artifacts_for_commit(&db, 7, "abc").await.unwrap();
607 assert!(artifacts_for_repo(&db, 7).await.unwrap().is_empty());
608 }
609}