anvilsign in

collin/anvil

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