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!("{PIPELINE_PATH}: `image` is required")));
55 }
56 Ok(pipeline)
57}
58
59/// Create a queued CI run for a pushed commit.
60pub async fn enqueue(
61 db: &toasty::Db,
62 repo_id: i64,
63 commit: &str,
64 ref_name: &str,
65) -> 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 mut runs = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
95 .exec(&mut conn)
96 .await?;
97 runs.sort_by(|a, b| b.id.cmp(&a.id));
98 runs.truncate(limit);
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 mut runs = CiRun::filter(CiRun::fields().repo_id().eq(repo_id))
110 .exec(&mut conn)
111 .await?;
112 runs.retain(|r| r.commit == commit);
113 runs.sort_by(|a, b| b.id.cmp(&a.id));
114 Ok(runs.into_iter().next())
115}
116
117/// Mark a run as started (running).
118pub async fn mark_running(db: &toasty::Db, id: i64) -> Result<()> {
119 let Some(mut run) = get(db, id).await? else {
120 return Ok(());
121 };
122 let mut conn = db.clone();
123 run.update()
124 .status(status::RUNNING)
125 .started_at(crate::now())
126 .exec(&mut conn)
127 .await?;
128 Ok(())
129}
130
131/// Append a chunk to a run's log.
132pub async fn append_log(db: &toasty::Db, id: i64, chunk: &str) -> Result<()> {
133 let Some(mut run) = get(db, id).await? else {
134 return Ok(());
135 };
136 let combined = format!("{}{chunk}", run.log);
137 let mut conn = db.clone();
138 run.update().log(&combined).exec(&mut conn).await?;
139 Ok(())
140}
141
142/// Finish a run with a terminal status (`success`/`failure`/`error`).
143pub async fn finish(db: &toasty::Db, id: i64, status: &str) -> Result<()> {
144 let Some(mut run) = get(db, id).await? else {
145 return Ok(());
146 };
147 let mut conn = db.clone();
148 run.update()
149 .status(status)
150 .finished_at(crate::now())
151 .exec(&mut conn)
152 .await?;
153 Ok(())
154}
155
156/// Re-queue runs left mid-flight by a crash/restart (status `running`).
157/// Returns the ids that were requeued so the runner can pick them up.
158pub async fn requeue_interrupted(db: &toasty::Db) -> Result<Vec<i64>> {
159 let mut conn = db.clone();
160 let interrupted = CiRun::filter(CiRun::fields().status().eq(status::RUNNING))
161 .exec(&mut conn)
162 .await?;
163 let mut ids = Vec::new();
164 for mut run in interrupted {
165 let mut conn = db.clone();
166 run.update().status(status::QUEUED).exec(&mut conn).await?;
167 ids.push(run.id);
168 }
169 Ok(ids)
170}
171
172/// List all queued run ids (oldest first) — used on startup to drain the queue.
173pub async fn queued_ids(db: &toasty::Db) -> Result<Vec<i64>> {
174 let mut conn = db.clone();
175 let mut runs = CiRun::filter(CiRun::fields().status().eq(status::QUEUED))
176 .exec(&mut conn)
177 .await?;
178 runs.sort_by(|a, b| a.id.cmp(&b.id));
179 Ok(runs.into_iter().map(|r| r.id).collect())
180}
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185
186 #[test]
187 fn parses_a_basic_pipeline() {
188 // YAML is indentation-sensitive, so the fixture is flush-left.
189 let p = parse_pipeline(
190 r#"image: rust:1.95-bookworm
191steps:
192 - name: test
193 run: cargo test --workspace
194 - run: cargo build --release
195"#,
196 )
197 .unwrap();
198 assert_eq!(p.image, "rust:1.95-bookworm");
199 assert_eq!(p.steps.len(), 2);
200 assert_eq!(p.steps[0].label(), "test");
201 // Unnamed step falls back to its command for the label.
202 assert_eq!(p.steps[1].label(), "cargo build --release");
203 }
204
205 #[test]
206 fn requires_an_image() {
207 assert!(parse_pipeline("steps: []\n").is_err());
208 }
209}