| 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 | |
| 12 | use std::{ |
| 13 | collections::BTreeMap, |
| 14 | io::Read, |
| 15 | }; |
| 16 | |
| 17 | use 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 | }; |
| 28 | use 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 | }; |
| 45 | use 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] |
| 60 | pub 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. |
| 75 | pub 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. |
| 211 | async 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). |
| 267 | async 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. |
| 292 | fn 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)] |
| 324 | mod 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 | } |