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::error::{Error, Result};
7use crate::models::CiRun;
8
9/// Run status values stored in [`CiRun::status`].
10pub mod status {
11 pub const QUEUED: &str = "queued";
12 pub const RUNNING: &str = "running";
13 pub const SUCCESS: &str = "success";
14 pub const FAILURE: &str = "failure";
15 pub const ERROR: &str = "error";
16}
17
18/// Path of the pipeline definition within a repository.
19pub const PIPELINE_PATH: &str = ".anvil/ci.yml";
20
21/// A parsed pipeline: a base image and ordered straight-line steps.
22#[derive(Debug, Clone, Deserialize)]
23pub struct Pipeline {
24 /// Docker image the steps run in, e.g. `rust:1.95-bookworm`.
25 pub image: String,
26 #[serde(default)]
27 pub steps: Vec<Step>,
28}
29
30/// One pipeline step: a shell command, with an optional display name.
31#[derive(Debug, Clone, Deserialize)]
32pub struct Step {
33 #[serde(default)]
34 pub name: String,
35 pub run: String,
36}
37
38impl Step {
39 /// Display label: the explicit name, or the command if unnamed.
40 pub fn label(&self) -> &str {
41 if self.name.is_empty() {
42 &self.run
43 } else {
44 &self.name
45 }
46 }
47}
48
49/// Parse a `.anvil/ci.yml` pipeline definition.
50pub fn parse_pipeline(yaml: &str) -> Result<Pipeline> {
51 let pipeline: Pipeline = serde_yaml::from_str(yaml)
52 .map_err(|e| Error::Invalid(format!("invalid {PIPELINE_PATH}: {e}")))?;
53 if pipeline.image.trim().is_empty() {
54 return Err(Error::Invalid(format!(
55 "{PIPELINE_PATH}: `image` is required"
56 )));
57 }
58 Ok(pipeline)
59}
60
61/// Create a queued CI run for a pushed commit.
62pub async fn enqueue(db: &toasty::Db, repo_id: i64, commit: &str, ref_name: &str) -> Result<CiRun> {
63 let mut conn = db.clone();
64 let run = toasty::create!(CiRun {
65 repo_id: repo_id,
66 commit: commit,
67 ref_name: ref_name,
68 status: status::QUEUED,
69 log: "",
70 created_at: crate::now(),
71 started_at: 0,
72 finished_at: 0,
73 })
74 .exec(&mut conn)
75 .await?;
76 Ok(run)
77}
78
79/// Fetch a run by id.
80pub async fn get(db: &toasty::Db, id: i64) -> Result<Option<CiRun>> {
81 let mut conn = db.clone();
82 Ok(CiRun::filter(CiRun::fields().id().eq(id))
83 .first()
84 .exec(&mut conn)
85 .await?)
86}
87
88/// List a repository's runs, newest first, up to `limit`.
89pub async fn list_by_repo(db: &toasty::Db, repo_id: i64, limit: usize) -> Result<Vec<CiRun>> {
90 let mut conn = db.clone();
91 let runs = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
92 .order_by(CiRun::fields().id().desc())
93 .limit(limit)
94 .exec(&mut conn)
95 .await?;
96 Ok(runs)
97}
98
99/// The most recent run for a specific commit (for status badges).
100pub async fn latest_for_commit(
101 db: &toasty::Db,
102 repo_id: i64,
103 commit: &str,
104) -> Result<Option<CiRun>> {
105 let mut conn = db.clone();
106 let run = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
107 .filter(CiRun::fields().commit().eq(commit))
108 .order_by(CiRun::fields().id().desc())
109 .first()
110 .exec(&mut conn)
111 .await?;
112 Ok(run)
113}
114
115/// Mark a run as started (running).
116pub async fn mark_running(db: &toasty::Db, id: i64) -> Result<()> {
117 let Some(mut run) = get(db, id).await? else {
118 return Ok(());
119 };
120 let mut conn = db.clone();
121 run.update()
122 .status(status::RUNNING)
123 .started_at(crate::now())
124 .exec(&mut conn)
125 .await?;
126 Ok(())
127}
128
129/// Append a chunk to a run's log.
130pub async fn append_log(db: &toasty::Db, id: i64, chunk: &str) -> Result<()> {
131 let Some(mut run) = get(db, id).await? else {
132 return Ok(());
133 };
134 let combined = format!("{}{chunk}", run.log);
135 let mut conn = db.clone();
136 run.update().log(&combined).exec(&mut conn).await?;
137 Ok(())
138}
139
140/// Finish a run with a terminal status (`success`/`failure`/`error`).
141pub async fn finish(db: &toasty::Db, id: i64, status: &str) -> Result<()> {
142 let Some(mut run) = get(db, id).await? else {
143 return Ok(());
144 };
145 let mut conn = db.clone();
146 run.update()
147 .status(status)
148 .finished_at(crate::now())
149 .exec(&mut conn)
150 .await?;
151 Ok(())
152}
153
154/// Re-queue runs left mid-flight by a crash/restart (status `running`).
155/// Returns the ids that were requeued so the runner can pick them up.
156pub async fn requeue_interrupted(db: &toasty::Db) -> Result<Vec<i64>> {
157 let mut conn = db.clone();
158 let interrupted = CiRun::filter(CiRun::fields().status().eq(status::RUNNING))
159 .exec(&mut conn)
160 .await?;
161 let mut ids = Vec::new();
162 for mut run in interrupted {
163 let mut conn = db.clone();
164 run.update().status(status::QUEUED).exec(&mut conn).await?;
165 ids.push(run.id);
166 }
167 Ok(ids)
168}
169
170/// List all queued run ids (oldest first) — used on startup to drain the queue.
171pub async fn queued_ids(db: &toasty::Db) -> Result<Vec<i64>> {
172 let mut conn = db.clone();
173 let runs = CiRun::filter(CiRun::fields().status().eq(status::QUEUED))
174 .order_by(CiRun::fields().id().asc())
175 .exec(&mut conn)
176 .await?;
177 Ok(runs.into_iter().map(|r| r.id).collect())
178}
179
180#[cfg(test)]
181mod tests {
182 use super::*;
183
184 #[test]
185 fn parses_a_basic_pipeline() {
186 // YAML is indentation-sensitive, so the fixture is flush-left.
187 let p = parse_pipeline(
188 r#"image: rust:1.95-bookworm
189steps:
190 - name: test
191 run: cargo test --workspace
192 - run: cargo build --release
193"#,
194 )
195 .unwrap();
196 assert_eq!(p.image, "rust:1.95-bookworm");
197 assert_eq!(p.steps.len(), 2);
198 assert_eq!(p.steps[0].label(), "test");
199 // Unnamed step falls back to its command for the label.
200 assert_eq!(p.steps[1].label(), "cargo build --release");
201 }
202
203 #[test]
204 fn requires_an_image() {
205 assert!(parse_pipeline("steps: []\n").is_err());
206 }
207
208 // Exercises the run-lifecycle queries against a real SQLite database, to
209 // confirm the ORM-level `order_by`/`limit`/filter actually work (Toasty is
210 // pre-1.0, so we don't take that on faith).
211 #[tokio::test]
212 async fn ordering_and_limit_run_in_the_database() {
213 let dir = tempfile::tempdir().unwrap();
214 let db = crate::db::connect(dir.path().join("t.db")).await.unwrap();
215
216 for c in ["aaa", "bbb", "ccc"] {
217 enqueue(&db, 1, c, "main").await.unwrap();
218 }
219 enqueue(&db, 2, "zzz", "main").await.unwrap(); // a different repo
220
221 // Newest-first, limited to 2 — and scoped to repo 1.
222 let runs = list_by_repo(&db, 1, 2).await.unwrap();
223 assert_eq!(runs.len(), 2);
224 assert!(runs[0].id > runs[1].id, "newest first");
225 assert_eq!(runs[0].commit, "ccc");
226 assert!(runs.iter().all(|r| r.repo_id == 1));
227
228 // Commit filter pushed into the query (not an in-memory retain).
229 let latest = latest_for_commit(&db, 1, "bbb").await.unwrap().unwrap();
230 assert_eq!(latest.commit, "bbb");
231 assert_eq!(latest.repo_id, 1);
232 assert!(latest_for_commit(&db, 1, "nope").await.unwrap().is_none());
233
234 // All queued, ascending by id.
235 let queued = queued_ids(&db).await.unwrap();
236 assert_eq!(queued.len(), 4);
237 assert!(queued.windows(2).all(|w| w[0] < w[1]), "ascending");
238 }
239}