| 1 | //! CI runner: drains queued [`anvil_core::ci`] runs and executes their |
| 2 | //! `.anvil/ci.yml` pipeline in a Docker container (via the socket, using |
| 3 | //! bollard). |
| 4 | //! |
| 5 | //! For each run: resolve the repo, materialize the commit's tree, parse the |
| 6 | //! pipeline, then run all steps as one `set -e` shell script inside the |
| 7 | //! pipeline's image. The checkout is uploaded into the container as a tar (via |
| 8 | //! the Docker API), so it works regardless of where anvil's own filesystem |
| 9 | //! lives and never exposes anvil's data volume to CI. |
| 10 | |
| 11 | use anvil_core::ci::{self, Pipeline}; |
| 12 | use anvil_core::{App, repos, storage, users}; |
| 13 | use anvil_git::browse::{self, TreeFile}; |
| 14 | use bollard::Docker; |
| 15 | use bollard::container::{ |
| 16 | Config, CreateContainerOptions, LogsOptions, RemoveContainerOptions, StartContainerOptions, |
| 17 | UploadToContainerOptions, WaitContainerOptions, |
| 18 | }; |
| 19 | use bollard::image::CreateImageOptions; |
| 20 | use futures_util::StreamExt; |
| 21 | use tokio::sync::mpsc::UnboundedReceiver; |
| 22 | |
| 23 | const WORKDIR: &str = "/workspace"; |
| 24 | |
| 25 | /// Run the CI worker loop: recover interrupted runs, drain the queue, then |
| 26 | /// process run ids as they arrive on `rx`. Runs one job at a time. |
| 27 | pub async fn run_worker(app: App, mut rx: UnboundedReceiver<i64>) { |
| 28 | match ci::requeue_interrupted(&app.db).await { |
| 29 | Ok(ids) if !ids.is_empty() => { |
| 30 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) |
| 31 | } |
| 32 | Ok(_) => {} |
| 33 | Err(e) => tracing::error!("ci: requeue failed: {e}"), |
| 34 | } |
| 35 | match ci::queued_ids(&app.db).await { |
| 36 | Ok(ids) => { |
| 37 | for id in ids { |
| 38 | run_one(&app, id).await; |
| 39 | } |
| 40 | } |
| 41 | Err(e) => tracing::error!("ci: listing queued runs failed: {e}"), |
| 42 | } |
| 43 | tracing::info!("ci runner ready"); |
| 44 | while let Some(id) = rx.recv().await { |
| 45 | run_one(&app, id).await; |
| 46 | } |
| 47 | } |
| 48 | |
| 49 | async fn run_one(app: &App, run_id: i64) { |
| 50 | tracing::info!("ci: run {run_id} starting"); |
| 51 | if let Err(e) = process(app, run_id).await { |
| 52 | tracing::error!("ci: run {run_id} errored: {e}"); |
| 53 | let _ = ci::append_log(&app.db, run_id, &format!("\n[runner error] {e}\n")).await; |
| 54 | let _ = ci::finish(&app.db, run_id, ci::status::ERROR).await; |
| 55 | } |
| 56 | } |
| 57 | |
| 58 | async fn process(app: &App, run_id: i64) -> Result<(), String> { |
| 59 | let run = ci::get(&app.db, run_id) |
| 60 | .await |
| 61 | .map_err(|e| e.to_string())? |
| 62 | .ok_or("run not found")?; |
| 63 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 64 | .await |
| 65 | .map_err(|e| e.to_string())? |
| 66 | .ok_or("repository not found")?; |
| 67 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 68 | .await |
| 69 | .map_err(|e| e.to_string())? |
| 70 | .ok_or("owner not found")?; |
| 71 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 72 | |
| 73 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 74 | .map_err(|e| e.to_string())? |
| 75 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 76 | let pipeline = |
| 77 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 78 | |
| 79 | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; |
| 80 | let tar = build_tar(&files); |
| 81 | |
| 82 | ci::mark_running(&app.db, run_id).await.ok(); |
| 83 | |
| 84 | let short = &run.commit[..run.commit.len().min(12)]; |
| 85 | let mut log = format!( |
| 86 | "anvil ci · {}/{} · {} @ {short}\nimage: {}\n", |
| 87 | owner.username, repo.name, run.ref_name, pipeline.image |
| 88 | ); |
| 89 | |
| 90 | let status = match execute(&pipeline, tar, &mut log).await { |
| 91 | Ok(0) => ci::status::SUCCESS, |
| 92 | Ok(code) => { |
| 93 | log.push_str(&format!("\n[exited with status {code}]\n")); |
| 94 | ci::status::FAILURE |
| 95 | } |
| 96 | Err(e) => { |
| 97 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 98 | ci::status::ERROR |
| 99 | } |
| 100 | }; |
| 101 | |
| 102 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 103 | ci::finish(&app.db, run_id, status).await.ok(); |
| 104 | tracing::info!("ci: run {run_id} {status}"); |
| 105 | |
| 106 | // Continuous deployment: on a green run of the configured deploy repo's |
| 107 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — |
| 108 | // no other repository can trigger it, even with passing CI. |
| 109 | if status == ci::status::SUCCESS |
| 110 | && app |
| 111 | .config |
| 112 | .ci |
| 113 | .is_deploy_target(&owner.username, &repo.name, &run.ref_name) |
| 114 | { |
| 115 | deploy(app, &owner.username, &repo.name, &run).await; |
| 116 | } |
| 117 | Ok(()) |
| 118 | } |
| 119 | |
| 120 | /// POST the configured deploy webhook. Best-effort: logs success/failure but |
| 121 | /// never fails the run (CI already passed). |
| 122 | async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) { |
| 123 | let cfg = &app.config.ci; |
| 124 | let body = serde_json::json!({ |
| 125 | "repo": format!("{owner}/{name}"), |
| 126 | "ref": run.ref_name, |
| 127 | "commit": run.commit, |
| 128 | "run_id": run.id, |
| 129 | }); |
| 130 | let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body); |
| 131 | if !cfg.deploy_secret.is_empty() { |
| 132 | req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret); |
| 133 | } |
| 134 | match req.send().await { |
| 135 | Ok(resp) if resp.status().is_success() => { |
| 136 | tracing::info!( |
| 137 | "ci: deploy webhook for {owner}/{name} accepted ({})", |
| 138 | resp.status() |
| 139 | ) |
| 140 | } |
| 141 | Ok(resp) => tracing::error!( |
| 142 | "ci: deploy webhook for {owner}/{name} returned {}", |
| 143 | resp.status() |
| 144 | ), |
| 145 | Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"), |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | /// Execute the pipeline in a container, streaming output into `log`. Returns the |
| 150 | /// container's exit code. |
| 151 | async fn execute(pipeline: &Pipeline, tar: Vec<u8>, log: &mut String) -> Result<i64, String> { |
| 152 | let docker = Docker::connect_with_socket_defaults() |
| 153 | .map_err(|e| format!("docker unavailable (is the socket mounted?): {e}"))?; |
| 154 | |
| 155 | // Pull the image (split name:tag so we don't accidentally pull all tags). |
| 156 | let (from_image, tag) = match pipeline.image.rsplit_once(':') { |
| 157 | Some((name, tag)) if !tag.contains('/') => (name.to_string(), tag.to_string()), |
| 158 | _ => (pipeline.image.clone(), "latest".to_string()), |
| 159 | }; |
| 160 | let mut pull = docker.create_image( |
| 161 | Some(CreateImageOptions { |
| 162 | from_image, |
| 163 | tag, |
| 164 | ..Default::default() |
| 165 | }), |
| 166 | None, |
| 167 | None, |
| 168 | ); |
| 169 | while let Some(item) = pull.next().await { |
| 170 | item.map_err(|e| format!("pull {}: {e}", pipeline.image))?; |
| 171 | } |
| 172 | |
| 173 | // Build a single `set -e` script from the steps. |
| 174 | let mut script = String::from("set -e\n"); |
| 175 | for step in &pipeline.steps { |
| 176 | script.push_str("printf '\\n=== %s ===\\n' "); |
| 177 | script.push_str(&single_quote(step.label())); |
| 178 | script.push('\n'); |
| 179 | script.push_str(&step.run); |
| 180 | script.push('\n'); |
| 181 | } |
| 182 | |
| 183 | let config = Config { |
| 184 | image: Some(pipeline.image.clone()), |
| 185 | cmd: Some(vec!["sh".to_string(), "-c".to_string(), script]), |
| 186 | working_dir: Some(WORKDIR.to_string()), |
| 187 | ..Default::default() |
| 188 | }; |
| 189 | let created = docker |
| 190 | .create_container(None::<CreateContainerOptions<String>>, config) |
| 191 | .await |
| 192 | .map_err(|e| format!("create container: {e}"))?; |
| 193 | let id = created.id; |
| 194 | |
| 195 | // Upload the checkout (tar entries are under `workspace/`, extracted at `/`). |
| 196 | docker |
| 197 | .upload_to_container( |
| 198 | &id, |
| 199 | Some(UploadToContainerOptions { |
| 200 | path: "/".to_string(), |
| 201 | ..Default::default() |
| 202 | }), |
| 203 | tar.into(), |
| 204 | ) |
| 205 | .await |
| 206 | .map_err(|e| format!("upload checkout: {e}"))?; |
| 207 | |
| 208 | docker |
| 209 | .start_container(&id, None::<StartContainerOptions<String>>) |
| 210 | .await |
| 211 | .map_err(|e| format!("start container: {e}"))?; |
| 212 | |
| 213 | // Stream logs until the container stops. |
| 214 | let mut logs = docker.logs( |
| 215 | &id, |
| 216 | Some(LogsOptions::<String> { |
| 217 | follow: true, |
| 218 | stdout: true, |
| 219 | stderr: true, |
| 220 | ..Default::default() |
| 221 | }), |
| 222 | ); |
| 223 | while let Some(item) = logs.next().await { |
| 224 | match item { |
| 225 | Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())), |
| 226 | Err(e) => { |
| 227 | log.push_str(&format!("\n[log stream error] {e}\n")); |
| 228 | break; |
| 229 | } |
| 230 | } |
| 231 | } |
| 232 | |
| 233 | // Wait for the exit code (non-zero surfaces as a wait error in bollard). |
| 234 | let mut code = 0i64; |
| 235 | let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>); |
| 236 | while let Some(item) = wait.next().await { |
| 237 | match item { |
| 238 | Ok(resp) => code = resp.status_code, |
| 239 | Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c, |
| 240 | Err(e) => { |
| 241 | let _ = docker |
| 242 | .remove_container( |
| 243 | &id, |
| 244 | Some(RemoveContainerOptions { |
| 245 | force: true, |
| 246 | ..Default::default() |
| 247 | }), |
| 248 | ) |
| 249 | .await; |
| 250 | return Err(format!("wait: {e}")); |
| 251 | } |
| 252 | } |
| 253 | } |
| 254 | |
| 255 | let _ = docker |
| 256 | .remove_container( |
| 257 | &id, |
| 258 | Some(RemoveContainerOptions { |
| 259 | force: true, |
| 260 | ..Default::default() |
| 261 | }), |
| 262 | ) |
| 263 | .await; |
| 264 | |
| 265 | Ok(code) |
| 266 | } |
| 267 | |
| 268 | /// Build an uncompressed tar of the checkout, rooted at `workspace/` so it |
| 269 | /// extracts to `/workspace` when uploaded to the container root. |
| 270 | fn build_tar(files: &[TreeFile]) -> Vec<u8> { |
| 271 | let mut builder = tar::Builder::new(Vec::new()); |
| 272 | for f in files { |
| 273 | let mut header = tar::Header::new_gnu(); |
| 274 | header.set_size(f.content.len() as u64); |
| 275 | header.set_mode(if f.executable { 0o755 } else { 0o644 }); |
| 276 | // append_data sets the path and checksum. |
| 277 | let _ = builder.append_data( |
| 278 | &mut header, |
| 279 | format!("workspace/{}", f.path), |
| 280 | f.content.as_slice(), |
| 281 | ); |
| 282 | } |
| 283 | builder.into_inner().unwrap_or_default() |
| 284 | } |
| 285 | |
| 286 | /// Single-quote a string for safe interpolation into a shell command. |
| 287 | fn single_quote(s: &str) -> String { |
| 288 | format!("'{}'", s.replace('\'', "'\\''")) |
| 289 | } |