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/// List all queued run ids (oldest first) — used on startup to drain the queue.
260pub async fn queued_ids(db: &toasty::Db) -> Result<Vec<i64>> {
261 let mut conn = db.clone();
262 let runs = CiRun::filter(CiRun::fields().status().eq(status::QUEUED))
263 .order_by(CiRun::fields().id().asc())
264 .exec(&mut conn)
265 .await?;
266 Ok(runs.into_iter().map(|r| r.id).collect())
267}
268
269/// Record a collected artifact. `meta` is a JSON object string (`{}` if none).
270#[allow(clippy::too_many_arguments)]
271pub async fn add_artifact(
272 db: &toasty::Db,
273 run_id: i64,
274 repo_id: i64,
275 commit: &str,
276 name: &str,
277 size: i64,
278 is_dir: bool,
279 browse: bool,
280 meta: &str,
281) -> Result<CiArtifact> {
282 let mut conn = db.clone();
283 let artifact = toasty::create!(CiArtifact {
284 run_id: run_id,
285 repo_id: repo_id,
286 commit: commit,
287 name: name,
288 size: size,
289 is_dir: is_dir,
290 browse: browse,
291 meta: meta,
292 created_at: crate::now(),
293 })
294 .exec(&mut conn)
295 .await?;
296 Ok(artifact)
297}
298
299/// A run's artifacts, in declaration (insertion) order.
300pub async fn artifacts_for_run(db: &toasty::Db, run_id: i64) -> Result<Vec<CiArtifact>> {
301 let mut conn = db.clone();
302 let artifacts = CiArtifact::filter(CiArtifact::fields().run_id().eq(run_id))
303 .order_by(CiArtifact::fields().id().asc())
304 .exec(&mut conn)
305 .await?;
306 Ok(artifacts)
307}
308
309/// The newest artifact named `name` for `commit` (across that commit's runs).
310pub async fn latest_artifact(
311 db: &toasty::Db,
312 repo_id: i64,
313 commit: &str,
314 name: &str,
315) -> Result<Option<CiArtifact>> {
316 let mut conn = db.clone();
317 let artifact = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
318 .filter(CiArtifact::fields().commit().eq(commit))
319 .filter(CiArtifact::fields().name().eq(name))
320 .order_by(CiArtifact::fields().id().desc())
321 .first()
322 .exec(&mut conn)
323 .await?;
324 Ok(artifact)
325}
326
327/// All artifact rows for a repository (for GC accounting). Small at our scale.
328pub async fn artifacts_for_repo(db: &toasty::Db, repo_id: i64) -> Result<Vec<CiArtifact>> {
329 let mut conn = db.clone();
330 let artifacts = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
331 .exec(&mut conn)
332 .await?;
333 Ok(artifacts)
334}
335
336/// Delete the artifact rows for one commit (the on-disk directory is the
337/// caller's to remove). Used when re-running a commit and by GC.
338pub async fn delete_artifacts_for_commit(
339 db: &toasty::Db,
340 repo_id: i64,
341 commit: &str,
342) -> Result<()> {
343 let mut conn = db.clone();
344 let rows = CiArtifact::filter(CiArtifact::fields().repo_id().eq(repo_id))
345 .filter(CiArtifact::fields().commit().eq(commit))
346 .exec(&mut conn)
347 .await?;
348 for row in rows {
349 let mut conn = db.clone();
350 row.delete().exec(&mut conn).await?;
351 }
352 Ok(())
353}
354
355#[cfg(test)]
356mod tests {
357 use super::*;
358
359 #[test]
360 fn parses_a_basic_pipeline() {
361 // YAML is indentation-sensitive, so the fixture is flush-left.
362 let p = parse_pipeline(
363 r#"image: rust:1.95-bookworm
364steps:
365 - name: test
366 run: cargo test --workspace
367 - run: cargo build --release
368"#,
369 )
370 .unwrap();
371 assert_eq!(p.image, "rust:1.95-bookworm");
372 assert_eq!(p.steps.len(), 2);
373 assert_eq!(p.steps[0].label(), "test");
374 // Unnamed step falls back to its command for the label.
375 assert_eq!(p.steps[1].label(), "cargo build --release");
376 }
377
378 #[test]
379 fn requires_an_image() {
380 assert!(parse_pipeline("steps: []\n").is_err());
381 }
382
383 #[test]
384 fn parses_and_validates_artifacts() {
385 let p = parse_pipeline(
386 r#"image: rust:1.95
387artifacts:
388 - name: anvild
389 path: target/release/anvild
390 meta:
391 version: ./target/release/anvild --version
392 - name: doc
393 path: target/doc/
394 browse: true
395"#,
396 )
397 .unwrap();
398 assert_eq!(p.artifacts.len(), 2);
399 assert_eq!(p.artifacts[0].name, "anvild");
400 assert_eq!(
401 p.artifacts[0].meta["version"],
402 "./target/release/anvild --version"
403 );
404 assert!(!p.artifacts[0].browse);
405 assert_eq!(p.artifacts[1].path, "target/doc", "trailing slash trimmed");
406 assert!(p.artifacts[1].browse);
407
408 let must_fail = |yaml: &str, why: &str| {
409 assert!(
410 parse_pipeline(&format!("image: i\n{yaml}")).is_err(),
411 "{why}"
412 );
413 };
414 must_fail(
415 "artifacts: [{name: 'a b', path: x}]",
416 "space in name rejected",
417 );
418 must_fail(
419 "artifacts: [{name: 'a/b', path: x}]",
420 "slash in name rejected",
421 );
422 must_fail("artifacts: [{name: '', path: x}]", "empty name rejected");
423 must_fail(
424 "artifacts: [{name: a, path: x}, {name: a, path: y}]",
425 "duplicate name rejected",
426 );
427 must_fail(
428 "artifacts: [{name: a, path: /etc}]",
429 "absolute path rejected",
430 );
431 must_fail(
432 "artifacts: [{name: a, path: '../up'}]",
433 "traversal rejected",
434 );
435 must_fail(
436 "artifacts: [{name: a, path: 'x//y'}]",
437 "empty segment rejected",
438 );
439 must_fail("artifacts: [{name: a, path: ''}]", "empty path rejected");
440 }
441
442 // Exercises the run-lifecycle queries against a real SQLite database, to
443 // confirm the ORM-level `order_by`/`limit`/filter actually work (Toasty is
444 // pre-1.0, so we don't take that on faith).
445 #[tokio::test]
446 async fn ordering_and_limit_run_in_the_database() {
447 let dir = tempfile::tempdir().unwrap();
448 let db = crate::db::connect(dir.path().join("t.db")).await.unwrap();
449
450 for c in ["aaa", "bbb", "ccc"] {
451 enqueue(&db, 1, c, "main").await.unwrap();
452 }
453 enqueue(&db, 2, "zzz", "main").await.unwrap(); // a different repo
454
455 // Newest-first, limited to 2 — and scoped to repo 1.
456 let runs = list_by_repo(&db, 1, 2).await.unwrap();
457 assert_eq!(runs.len(), 2);
458 assert!(runs[0].id > runs[1].id, "newest first");
459 assert_eq!(runs[0].commit, "ccc");
460 assert!(runs.iter().all(|r| r.repo_id == 1));
461
462 // Commit filter pushed into the query (not an in-memory retain).
463 let latest = latest_for_commit(&db, 1, "bbb").await.unwrap().unwrap();
464 assert_eq!(latest.commit, "bbb");
465 assert_eq!(latest.repo_id, 1);
466 assert!(latest_for_commit(&db, 1, "nope").await.unwrap().is_none());
467
468 // All queued, ascending by id.
469 let queued = queued_ids(&db).await.unwrap();
470 assert_eq!(queued.len(), 4);
471 assert!(queued.windows(2).all(|w| w[0] < w[1]), "ascending");
472 }
473
474 #[tokio::test]
475 async fn artifact_rows_round_trip() {
476 let dir = tempfile::tempdir().unwrap();
477 let db = crate::db::connect(dir.path().join("t.db")).await.unwrap();
478
479 add_artifact(&db, 1, 7, "abc", "bin", 100, false, false, "{}")
480 .await
481 .unwrap();
482 add_artifact(&db, 1, 7, "abc", "doc", 5000, true, true, r#"{"v":"1"}"#)
483 .await
484 .unwrap();
485 // A newer run of the same commit supersedes for `latest_artifact`.
486 add_artifact(&db, 2, 7, "abc", "bin", 200, false, false, "{}")
487 .await
488 .unwrap();
489
490 let run1 = artifacts_for_run(&db, 1).await.unwrap();
491 assert_eq!(run1.len(), 2);
492 assert_eq!(run1[0].name, "bin");
493 assert!(run1[1].browse);
494
495 let latest = latest_artifact(&db, 7, "abc", "bin")
496 .await
497 .unwrap()
498 .unwrap();
499 assert_eq!(latest.run_id, 2);
500 assert_eq!(latest.size, 200);
501 assert!(
502 latest_artifact(&db, 8, "abc", "bin")
503 .await
504 .unwrap()
505 .is_none(),
506 "scoped to repo"
507 );
508
509 delete_artifacts_for_commit(&db, 7, "abc").await.unwrap();
510 assert!(artifacts_for_repo(&db, 7).await.unwrap().is_empty());
511 }
512}