anvilsign in

collin/anvil

main / crates / anvil-worker / src / executor.rs
1//! Running one job in a sandboxed container.
2//!
3//! Lifted verbatim out of `anvil-ci`, where it ran in anvild's own process
4//! against the host's Docker socket. The sandbox is unchanged — that is the
5//! point: moving execution to another machine must not quietly relax it.
6//!
7//! The job container never sees the Docker socket and gets no mounts of any
8//! kind. The checkout is *uploaded* through the Docker API and artifacts are
9//! *downloaded* back out the same way, so a job cannot reach the runner's
10//! filesystem any more than it could reach anvil's.
11
12use std::{
13 collections::BTreeMap,
14 io::Read,
15};
16
17use anvil_job::{
18 ArtifactSpec,
19 CollectedArtifact,
20 JobSpec,
21 META_DIR,
22 META_TAR_CAP,
23 META_VALUE_CAP,
24 Stored,
25 WORKDIR,
26 mb_cap,
27};
28use bollard::{
29 Docker,
30 body_full,
31 models::{
32 ContainerCreateBody,
33 HostConfig,
34 },
35 query_parameters::{
36 CreateContainerOptions,
37 DownloadFromContainerOptions,
38 LogsOptions,
39 RemoveContainerOptions,
40 StartContainerOptions,
41 UploadToContainerOptions,
42 WaitContainerOptions,
43 },
44};
45use futures_util::StreamExt;
46
47/// Where a collected artifact's bytes go.
48///
49/// A trait rather than a direct upload call so the download loop can hand each
50/// tar off and drop it rather than buffering every artifact —
51/// `artifact_run_max_mb` defaults to 512, which is not an amount to hold in
52/// RAM. Returning [`Stored`] rather than `()` is what lets the caller charge
53/// the run budget the size that actually landed on anvil's disk, which the
54/// runner cannot compute itself: a `browse` directory is extracted there and
55/// any other directory recompressed.
56// async-trait rewrites the method to return a boxed future, which is already
57// `#[must_use]`; the attribute it also emits then trips `double_must_use`.
58#[allow(clippy::double_must_use)]
59#[async_trait::async_trait]
60pub trait ArtifactSink: Send {
61 async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String>;
62}
63
64/// Execute a job in a sandboxed container, streaming output into `log`.
65/// Returns the exit code and whatever artifacts `sink` accepted (none on
66/// timeout — the container is already gone).
67///
68/// The job container never sees the Docker socket and gets no mounts of any
69/// kind (the checkout is *uploaded*, not bind-mounted; artifacts are
70/// *downloaded* out the same way). All capabilities are dropped and
71/// `no-new-privileges` is set unconditionally; pids/memory/cpu caps, the
72/// wall-clock timeout, network access and the container user come from
73/// `spec.sandbox`. The image allowlist was applied when the spec was built —
74/// a runner receiving a spec does not get to widen it.
75pub async fn execute(
76 spec: &JobSpec,
77 tar: Vec<u8>,
78 log: &mut String,
79 sink: &mut dyn ArtifactSink,
80) -> Result<(i64, Vec<CollectedArtifact>), String> {
81 let platform = spec.platform.clone().unwrap_or_default();
82 let docker = anvil_docker::connect()?;
83 anvil_docker::ensure_image(&docker, &spec.image, &platform).await?;
84 let sb = &spec.sandbox;
85
86 // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap.
87 let host_config = HostConfig {
88 cap_drop: Some(vec!["ALL".to_string()]),
89 security_opt: Some(vec!["no-new-privileges:true".to_string()]),
90 pids_limit: (sb.pids_limit > 0).then_some(sb.pids_limit),
91 memory: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024),
92 memory_swap: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024),
93 nano_cpus: (sb.cpus > 0.0).then_some((sb.cpus * 1e9) as i64),
94 network_mode: (!sb.network).then(|| "none".to_string()),
95 ..Default::default()
96 };
97 let config = ContainerCreateBody {
98 image: Some(spec.image.clone()),
99 cmd: Some(vec![
100 "sh".to_string(),
101 "-c".to_string(),
102 spec.script.clone(),
103 ]),
104 env: (!spec.env.is_empty())
105 .then(|| spec.env.iter().map(|(k, v)| format!("{k}={v}")).collect()),
106 working_dir: Some(WORKDIR.to_string()),
107 user: (!sb.run_as.is_empty()).then(|| sb.run_as.clone()),
108 host_config: Some(host_config),
109 ..Default::default()
110 };
111 // An explicit platform needs the options struct; without one, pass None so
112 // the daemon picks its native architecture exactly as before.
113 let opts = spec.platform.as_ref().map(|p| CreateContainerOptions {
114 name: None,
115 platform: p.clone(),
116 });
117 let created = docker
118 .create_container(opts, config)
119 .await
120 .map_err(|e| format!("create container: {e}"))?;
121 let id = created.id;
122
123 // Upload the checkout (tar entries are under `workspace/`, extracted at `/`).
124 docker
125 .upload_to_container(
126 &id,
127 Some(UploadToContainerOptions {
128 path: "/".to_string(),
129 ..Default::default()
130 }),
131 body_full(tar.into()),
132 )
133 .await
134 .map_err(|e| format!("upload checkout: {e}"))?;
135
136 docker
137 .start_container(&id, None::<StartContainerOptions>)
138 .await
139 .map_err(|e| format!("start container: {e}"))?;
140
141 // Stream logs and wait for the exit code, bounded by the wall-clock
142 // timeout. The container is force-removed on every path (which also kills
143 // a still-running job after a timeout).
144 let run = async {
145 let mut logs = docker.logs(
146 &id,
147 Some(LogsOptions {
148 follow: true,
149 stdout: true,
150 stderr: true,
151 ..Default::default()
152 }),
153 );
154 while let Some(item) = logs.next().await {
155 match item {
156 Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())),
157 Err(e) => {
158 log.push_str(&format!("\n[log stream error] {e}\n"));
159 break;
160 }
161 }
162 }
163
164 // Non-zero exit codes surface as a wait error in bollard.
165 let mut code = 0i64;
166 let mut wait = docker.wait_container(&id, None::<WaitContainerOptions>);
167 while let Some(item) = wait.next().await {
168 match item {
169 Ok(resp) => code = resp.status_code,
170 Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c,
171 Err(e) => return Err(format!("wait: {e}")),
172 }
173 }
174 Ok(code)
175 };
176 let result = match sb.timeout_secs {
177 0 => run.await,
178 secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run)
179 .await
180 .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))),
181 };
182
183 // Artifacts come out of the (now stopped) container before it is removed.
184 let collected = match &result {
185 Ok(_) if !spec.artifacts.is_empty() => {
186 collect_artifacts(&docker, &id, spec, log, sink).await
187 }
188 _ => Vec::new(),
189 };
190
191 let _ = docker
192 .remove_container(
193 &id,
194 Some(RemoveContainerOptions {
195 force: true,
196 ..Default::default()
197 }),
198 )
199 .await;
200
201 result.map(|code| (code, collected))
202}
203
204/// Collect the job's declared artifacts from the stopped container, handing
205/// each tar to `sink` as it is downloaded. Failures are per-artifact: each is
206/// logged and skipped, never failing the run.
207///
208/// One artifact is in memory at a time by construction — the tar is dropped
209/// once the sink has taken it — which is what keeps `artifact_run_max_mb`
210/// (512 MiB by default) a disk budget rather than a memory one.
211async fn collect_artifacts(
212 docker: &Docker,
213 id: &str,
214 job: &JobSpec,
215 log: &mut String,
216 sink: &mut dyn ArtifactSink,
217) -> Vec<CollectedArtifact> {
218 // Extractor outputs first: artifact name → key → value.
219 let mut metas: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new();
220 if job.artifacts.iter().any(|a| a.has_meta) {
221 match download_tar(docker, id, META_DIR, META_TAR_CAP).await {
222 Ok(Some(bytes)) => metas = parse_meta_tar(&bytes),
223 Ok(None) => log.push_str("\n[artifacts: extractor output exceeded its cap]\n"),
224 Err(e) => log.push_str(&format!("\n[artifacts: reading extractor output: {e}]\n")),
225 }
226 }
227
228 let per_artifact_cap = mb_cap(job.sandbox.artifact_max_mb);
229 let mut run_budget = mb_cap(job.sandbox.artifact_run_max_mb);
230 let mut collected = Vec::new();
231 for spec in &job.artifacts {
232 let cap = per_artifact_cap.min(run_budget);
233 let note = |log: &mut String, what: &str| {
234 log.push_str(&format!("\n[artifact {}: {what}]\n", spec.name));
235 };
236 let bytes = match download_tar(docker, id, &format!("{WORKDIR}/{}", spec.path), cap).await {
237 Ok(Some(bytes)) => bytes,
238 Ok(None) => {
239 note(log, "exceeds the size cap; skipped");
240 continue;
241 }
242 Err(e) => {
243 note(log, &format!("download failed ({e}); skipped"));
244 continue;
245 }
246 };
247 match sink.put(spec, &bytes).await {
248 Ok(Stored { size, is_dir }) => {
249 run_budget = run_budget.saturating_sub(size as u64);
250 let meta = metas.get(&spec.name).cloned().unwrap_or_default();
251 collected.push(CollectedArtifact {
252 name: spec.name.clone(),
253 size,
254 is_dir,
255 browse: spec.browse,
256 meta: serde_json::to_string(&meta).unwrap_or_else(|_| "{}".into()),
257 });
258 }
259 Err(e) => note(log, &format!("storing failed ({e}); skipped")),
260 }
261 }
262 collected
263}
264
265/// Download `path` from the container as a tar, buffering at most `cap` bytes
266/// (`Ok(None)` when exceeded).
267async fn download_tar(
268 docker: &Docker,
269 id: &str,
270 path: &str,
271 cap: u64,
272) -> Result<Option<Vec<u8>>, String> {
273 let mut stream = docker.download_from_container(
274 id,
275 Some(DownloadFromContainerOptions {
276 path: path.to_string(),
277 }),
278 );
279 let mut buf = Vec::new();
280 while let Some(chunk) = stream.next().await {
281 let chunk = chunk.map_err(|e| e.to_string())?;
282 if (buf.len() + chunk.len()) as u64 > cap {
283 return Ok(None);
284 }
285 buf.extend_from_slice(&chunk);
286 }
287 Ok(Some(buf))
288}
289
290/// Parse the extractor-output tar (`anvil-meta/<artifact>/<key>` files) into
291/// artifact → key → trimmed value.
292fn parse_meta_tar(bytes: &[u8]) -> BTreeMap<String, BTreeMap<String, String>> {
293 let mut out: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new();
294 let mut archive = tar::Archive::new(bytes);
295 let Ok(entries) = archive.entries() else {
296 return out;
297 };
298 for entry in entries.flatten() {
299 if !entry.header().entry_type().is_file() {
300 continue;
301 }
302 let Ok(path) = entry.path() else { continue };
303 // anvil-meta/<artifact>/<key>
304 let parts: Vec<String> = path
305 .components()
306 .skip(1)
307 .map(|c| c.as_os_str().to_string_lossy().into_owned())
308 .collect();
309 let [artifact, key] = parts.as_slice() else {
310 continue;
311 };
312 let (artifact, key) = (artifact.clone(), key.clone());
313 let mut value = String::new();
314 let _ = entry.take(META_VALUE_CAP as u64).read_to_string(&mut value);
315 let value = value.trim().to_string();
316 if !value.is_empty() {
317 out.entry(artifact).or_default().insert(key, value);
318 }
319 }
320 out
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326
327 /// A tar with the same layout the Docker archive endpoint produces:
328 /// entries rooted at the requested item's basename. `None` content means a
329 /// directory entry.
330 fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> {
331 let mut b = tar::Builder::new(Vec::new());
332 for (path, content) in entries {
333 let mut h = tar::Header::new_gnu();
334 match content {
335 Some(c) => {
336 h.set_size(c.len() as u64);
337 h.set_mode(0o644);
338 h.set_entry_type(tar::EntryType::Regular);
339 b.append_data(&mut h, path, c.as_bytes()).unwrap();
340 }
341 None => {
342 h.set_size(0);
343 h.set_mode(0o755);
344 h.set_entry_type(tar::EntryType::Directory);
345 b.append_data(&mut h, path, std::io::empty()).unwrap();
346 }
347 }
348 }
349 b.into_inner().unwrap()
350 }
351
352 #[test]
353 fn parses_meta_tar_with_trimmed_capped_values() {
354 let tar = tar_of(&[
355 ("anvil-meta", None),
356 ("anvil-meta/bin", None),
357 ("anvil-meta/bin/version", Some("anvild 0.0.0\n")),
358 ("anvil-meta/bin/empty", Some(" \n")),
359 ("anvil-meta/doc", None),
360 ("anvil-meta/doc/pages", Some("42")),
361 ]);
362 let metas = parse_meta_tar(&tar);
363 assert_eq!(metas["bin"]["version"], "anvild 0.0.0");
364 assert_eq!(metas["doc"]["pages"], "42");
365 assert!(!metas["bin"].contains_key("empty"), "blank values dropped");
366 }
367}