anvilsign in

collin/anvil

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