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