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