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