| 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 | container::{ |
| 31 | Config, |
| 32 | CreateContainerOptions, |
| 33 | DownloadFromContainerOptions, |
| 34 | LogsOptions, |
| 35 | RemoveContainerOptions, |
| 36 | StartContainerOptions, |
| 37 | UploadToContainerOptions, |
| 38 | WaitContainerOptions, |
| 39 | }, |
| 40 | models::HostConfig, |
| 41 | }; |
| 42 | use 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] |
| 57 | pub 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. |
| 72 | pub 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. |
| 208 | async 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). |
| 264 | async 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. |
| 289 | fn 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)] |
| 321 | mod 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 | } |