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::{
18 self,
19 Pipeline,
20};
21use anvil_core::config::CiConfig;
22use anvil_core::{
23 App,
24 repos,
25 storage,
26 users,
27};
28use anvil_git::browse::{
29 self,
30 TreeFile,
31};
32use bollard::Docker;
33use bollard::container::{
34 Config,
35 CreateContainerOptions,
36 LogsOptions,
37 RemoveContainerOptions,
38 StartContainerOptions,
39 UploadToContainerOptions,
40 WaitContainerOptions,
41};
42use bollard::image::CreateImageOptions;
43use bollard::models::HostConfig;
44use futures_util::StreamExt;
45use tokio::sync::mpsc::UnboundedReceiver;
46
47const 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.
51pub 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
73async 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
82async 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).
146async 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`.
181async 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.
324fn 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.
341fn single_quote(s: &str) -> String {
342 format!("'{}'", s.replace('\'', "'\\''"))
343}