| 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 | //! anvil is the *broker*: it is the only Docker client, and the job container |
| 12 | //! gets no socket, no bind mounts, and no volumes. On top of that the job runs |
| 13 | //! with all capabilities dropped, `no-new-privileges`, and configurable |
| 14 | //! pids/memory/cpu caps plus a wall-clock timeout and an optional image |
| 15 | //! allowlist ([`anvil_core::config::CiConfig`]). |
| 16 | |
| 17 | use anvil_core::ci::{ |
| 18 | self, |
| 19 | Pipeline, |
| 20 | }; |
| 21 | use anvil_core::config::CiConfig; |
| 22 | use anvil_core::{ |
| 23 | App, |
| 24 | repos, |
| 25 | storage, |
| 26 | users, |
| 27 | }; |
| 28 | use anvil_git::browse::{ |
| 29 | self, |
| 30 | TreeFile, |
| 31 | }; |
| 32 | use bollard::Docker; |
| 33 | use bollard::container::{ |
| 34 | Config, |
| 35 | CreateContainerOptions, |
| 36 | LogsOptions, |
| 37 | RemoveContainerOptions, |
| 38 | StartContainerOptions, |
| 39 | UploadToContainerOptions, |
| 40 | WaitContainerOptions, |
| 41 | }; |
| 42 | use bollard::image::CreateImageOptions; |
| 43 | use bollard::models::HostConfig; |
| 44 | use futures_util::StreamExt; |
| 45 | use tokio::sync::mpsc::UnboundedReceiver; |
| 46 | |
| 47 | const WORKDIR: &str = "/workspace"; |
| 48 | |
| 49 | /// Run the CI worker loop: recover interrupted runs, drain the queue, then |
| 50 | /// process run ids as they arrive on `rx`. Runs one job at a time. |
| 51 | pub async fn run_worker(app: App, mut rx: UnboundedReceiver<i64>) { |
| 52 | match ci::requeue_interrupted(&app.db).await { |
| 53 | Ok(ids) if !ids.is_empty() => { |
| 54 | tracing::info!("ci: requeued {} interrupted run(s)", ids.len()) |
| 55 | } |
| 56 | Ok(_) => {} |
| 57 | Err(e) => tracing::error!("ci: requeue failed: {e}"), |
| 58 | } |
| 59 | match ci::queued_ids(&app.db).await { |
| 60 | Ok(ids) => { |
| 61 | for id in ids { |
| 62 | run_one(&app, id).await; |
| 63 | } |
| 64 | } |
| 65 | Err(e) => tracing::error!("ci: listing queued runs failed: {e}"), |
| 66 | } |
| 67 | tracing::info!("ci runner ready"); |
| 68 | while let Some(id) = rx.recv().await { |
| 69 | run_one(&app, id).await; |
| 70 | } |
| 71 | } |
| 72 | |
| 73 | async fn run_one(app: &App, run_id: i64) { |
| 74 | tracing::info!("ci: run {run_id} starting"); |
| 75 | if let Err(e) = process(app, run_id).await { |
| 76 | tracing::error!("ci: run {run_id} errored: {e}"); |
| 77 | let _ = ci::append_log(&app.db, run_id, &format!("\n[runner error] {e}\n")).await; |
| 78 | let _ = ci::finish(&app.db, run_id, ci::status::ERROR).await; |
| 79 | } |
| 80 | } |
| 81 | |
| 82 | async fn process(app: &App, run_id: i64) -> Result<(), String> { |
| 83 | let run = ci::get(&app.db, run_id) |
| 84 | .await |
| 85 | .map_err(|e| e.to_string())? |
| 86 | .ok_or("run not found")?; |
| 87 | let repo = repos::find_by_id(&app.db, run.repo_id) |
| 88 | .await |
| 89 | .map_err(|e| e.to_string())? |
| 90 | .ok_or("repository not found")?; |
| 91 | let owner = users::find_by_id(&app.db, repo.owner_id) |
| 92 | .await |
| 93 | .map_err(|e| e.to_string())? |
| 94 | .ok_or("owner not found")?; |
| 95 | let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name); |
| 96 | |
| 97 | let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH) |
| 98 | .map_err(|e| e.to_string())? |
| 99 | .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?; |
| 100 | let pipeline = |
| 101 | ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?; |
| 102 | |
| 103 | let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?; |
| 104 | let tar = build_tar(&files); |
| 105 | |
| 106 | ci::mark_running(&app.db, run_id).await.ok(); |
| 107 | |
| 108 | let short = &run.commit[..run.commit.len().min(12)]; |
| 109 | let mut log = format!( |
| 110 | "anvil ci · {}/{} · {} @ {short}\nimage: {}\n", |
| 111 | owner.username, repo.name, run.ref_name, pipeline.image |
| 112 | ); |
| 113 | |
| 114 | let status = match execute(&pipeline, tar, &mut log, &app.config.ci).await { |
| 115 | Ok(0) => ci::status::SUCCESS, |
| 116 | Ok(code) => { |
| 117 | log.push_str(&format!("\n[exited with status {code}]\n")); |
| 118 | ci::status::FAILURE |
| 119 | } |
| 120 | Err(e) => { |
| 121 | log.push_str(&format!("\n[runner error] {e}\n")); |
| 122 | ci::status::ERROR |
| 123 | } |
| 124 | }; |
| 125 | |
| 126 | ci::append_log(&app.db, run_id, &log).await.ok(); |
| 127 | ci::finish(&app.db, run_id, status).await.ok(); |
| 128 | tracing::info!("ci: run {run_id} {status}"); |
| 129 | |
| 130 | // Continuous deployment: on a green run of the configured deploy repo's |
| 131 | // deploy branch, fire the redeploy webhook. Scoped to one repo by config — |
| 132 | // no other repository can trigger it, even with passing CI. |
| 133 | if status == ci::status::SUCCESS |
| 134 | && app |
| 135 | .config |
| 136 | .ci |
| 137 | .is_deploy_target(&owner.username, &repo.name, &run.ref_name) |
| 138 | { |
| 139 | deploy(app, &owner.username, &repo.name, &run).await; |
| 140 | } |
| 141 | Ok(()) |
| 142 | } |
| 143 | |
| 144 | /// POST the configured deploy webhook. Best-effort: logs success/failure but |
| 145 | /// never fails the run (CI already passed). |
| 146 | async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) { |
| 147 | let cfg = &app.config.ci; |
| 148 | let body = serde_json::json!({ |
| 149 | "repo": format!("{owner}/{name}"), |
| 150 | "ref": run.ref_name, |
| 151 | "commit": run.commit, |
| 152 | "run_id": run.id, |
| 153 | }); |
| 154 | let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body); |
| 155 | if !cfg.deploy_secret.is_empty() { |
| 156 | req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret); |
| 157 | } |
| 158 | match req.send().await { |
| 159 | Ok(resp) if resp.status().is_success() => { |
| 160 | tracing::info!( |
| 161 | "ci: deploy webhook for {owner}/{name} accepted ({})", |
| 162 | resp.status() |
| 163 | ) |
| 164 | } |
| 165 | Ok(resp) => tracing::error!( |
| 166 | "ci: deploy webhook for {owner}/{name} returned {}", |
| 167 | resp.status() |
| 168 | ), |
| 169 | Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"), |
| 170 | } |
| 171 | } |
| 172 | |
| 173 | /// Execute the pipeline in a sandboxed container, streaming output into `log`. |
| 174 | /// Returns the container's exit code. |
| 175 | /// |
| 176 | /// The job container never sees the Docker socket and gets no mounts of any |
| 177 | /// kind (the checkout is *uploaded*, not bind-mounted). All capabilities are |
| 178 | /// dropped and `no-new-privileges` is set unconditionally; pids/memory/cpu |
| 179 | /// caps, the wall-clock timeout, network access, the container user, and the |
| 180 | /// image allowlist come from `cfg`. |
| 181 | async fn execute( |
| 182 | pipeline: &Pipeline, |
| 183 | tar: Vec<u8>, |
| 184 | log: &mut String, |
| 185 | cfg: &CiConfig, |
| 186 | ) -> Result<i64, String> { |
| 187 | if !cfg.image_allowed(&pipeline.image) { |
| 188 | return Err(format!( |
| 189 | "image {} is not permitted by ci.allowed_images", |
| 190 | pipeline.image |
| 191 | )); |
| 192 | } |
| 193 | let docker = Docker::connect_with_socket_defaults() |
| 194 | .map_err(|e| format!("docker unavailable (is the socket mounted?): {e}"))?; |
| 195 | |
| 196 | // Pull the image (split name:tag so we don't accidentally pull all tags). |
| 197 | let (from_image, tag) = match pipeline.image.rsplit_once(':') { |
| 198 | Some((name, tag)) if !tag.contains('/') => (name.to_string(), tag.to_string()), |
| 199 | _ => (pipeline.image.clone(), "latest".to_string()), |
| 200 | }; |
| 201 | let mut pull = docker.create_image( |
| 202 | Some(CreateImageOptions { |
| 203 | from_image, |
| 204 | tag, |
| 205 | ..Default::default() |
| 206 | }), |
| 207 | None, |
| 208 | None, |
| 209 | ); |
| 210 | while let Some(item) = pull.next().await { |
| 211 | item.map_err(|e| format!("pull {}: {e}", pipeline.image))?; |
| 212 | } |
| 213 | |
| 214 | // Build a single `set -e` script from the steps. |
| 215 | let mut script = String::from("set -e\n"); |
| 216 | for step in &pipeline.steps { |
| 217 | script.push_str("printf '\\n=== %s ===\\n' "); |
| 218 | script.push_str(&single_quote(step.label())); |
| 219 | script.push('\n'); |
| 220 | script.push_str(&step.run); |
| 221 | script.push('\n'); |
| 222 | } |
| 223 | |
| 224 | // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap. |
| 225 | let host_config = HostConfig { |
| 226 | cap_drop: Some(vec!["ALL".to_string()]), |
| 227 | security_opt: Some(vec!["no-new-privileges:true".to_string()]), |
| 228 | pids_limit: (cfg.pids_limit > 0).then_some(cfg.pids_limit), |
| 229 | memory: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 230 | memory_swap: (cfg.memory_mb > 0).then(|| cfg.memory_mb * 1024 * 1024), |
| 231 | nano_cpus: (cfg.cpus > 0.0).then_some((cfg.cpus * 1e9) as i64), |
| 232 | network_mode: (!cfg.network).then(|| "none".to_string()), |
| 233 | ..Default::default() |
| 234 | }; |
| 235 | let config = Config { |
| 236 | image: Some(pipeline.image.clone()), |
| 237 | cmd: Some(vec!["sh".to_string(), "-c".to_string(), script]), |
| 238 | working_dir: Some(WORKDIR.to_string()), |
| 239 | user: (!cfg.run_as.is_empty()).then(|| cfg.run_as.clone()), |
| 240 | host_config: Some(host_config), |
| 241 | ..Default::default() |
| 242 | }; |
| 243 | let created = docker |
| 244 | .create_container(None::<CreateContainerOptions<String>>, config) |
| 245 | .await |
| 246 | .map_err(|e| format!("create container: {e}"))?; |
| 247 | let id = created.id; |
| 248 | |
| 249 | // Upload the checkout (tar entries are under `workspace/`, extracted at `/`). |
| 250 | docker |
| 251 | .upload_to_container( |
| 252 | &id, |
| 253 | Some(UploadToContainerOptions { |
| 254 | path: "/".to_string(), |
| 255 | ..Default::default() |
| 256 | }), |
| 257 | tar.into(), |
| 258 | ) |
| 259 | .await |
| 260 | .map_err(|e| format!("upload checkout: {e}"))?; |
| 261 | |
| 262 | docker |
| 263 | .start_container(&id, None::<StartContainerOptions<String>>) |
| 264 | .await |
| 265 | .map_err(|e| format!("start container: {e}"))?; |
| 266 | |
| 267 | // Stream logs and wait for the exit code, bounded by the wall-clock |
| 268 | // timeout. The container is force-removed on every path (which also kills |
| 269 | // a still-running job after a timeout). |
| 270 | let run = async { |
| 271 | let mut logs = docker.logs( |
| 272 | &id, |
| 273 | Some(LogsOptions::<String> { |
| 274 | follow: true, |
| 275 | stdout: true, |
| 276 | stderr: true, |
| 277 | ..Default::default() |
| 278 | }), |
| 279 | ); |
| 280 | while let Some(item) = logs.next().await { |
| 281 | match item { |
| 282 | Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())), |
| 283 | Err(e) => { |
| 284 | log.push_str(&format!("\n[log stream error] {e}\n")); |
| 285 | break; |
| 286 | } |
| 287 | } |
| 288 | } |
| 289 | |
| 290 | // Non-zero exit codes surface as a wait error in bollard. |
| 291 | let mut code = 0i64; |
| 292 | let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>); |
| 293 | while let Some(item) = wait.next().await { |
| 294 | match item { |
| 295 | Ok(resp) => code = resp.status_code, |
| 296 | Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c, |
| 297 | Err(e) => return Err(format!("wait: {e}")), |
| 298 | } |
| 299 | } |
| 300 | Ok(code) |
| 301 | }; |
| 302 | let result = match cfg.timeout_secs { |
| 303 | 0 => run.await, |
| 304 | secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run) |
| 305 | .await |
| 306 | .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))), |
| 307 | }; |
| 308 | |
| 309 | let _ = docker |
| 310 | .remove_container( |
| 311 | &id, |
| 312 | Some(RemoveContainerOptions { |
| 313 | force: true, |
| 314 | ..Default::default() |
| 315 | }), |
| 316 | ) |
| 317 | .await; |
| 318 | |
| 319 | result |
| 320 | } |
| 321 | |
| 322 | /// Build an uncompressed tar of the checkout, rooted at `workspace/` so it |
| 323 | /// extracts to `/workspace` when uploaded to the container root. |
| 324 | fn build_tar(files: &[TreeFile]) -> Vec<u8> { |
| 325 | let mut builder = tar::Builder::new(Vec::new()); |
| 326 | for f in files { |
| 327 | let mut header = tar::Header::new_gnu(); |
| 328 | header.set_size(f.content.len() as u64); |
| 329 | header.set_mode(if f.executable { 0o755 } else { 0o644 }); |
| 330 | // append_data sets the path and checksum. |
| 331 | let _ = builder.append_data( |
| 332 | &mut header, |
| 333 | format!("workspace/{}", f.path), |
| 334 | f.content.as_slice(), |
| 335 | ); |
| 336 | } |
| 337 | builder.into_inner().unwrap_or_default() |
| 338 | } |
| 339 | |
| 340 | /// Single-quote a string for safe interpolation into a shell command. |
| 341 | fn single_quote(s: &str) -> String { |
| 342 | format!("'{}'", s.replace('\'', "'\\''")) |
| 343 | } |