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