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 std::{
18 collections::BTreeMap,
19 io::Read,
20 path::Path,
21};
22
23use anvil_core::{
24 App,
25 ci::{
26 self,
27 Pipeline,
28 },
29 config::CiConfig,
30 repos,
31 storage,
32 users,
33};
34use anvil_git::browse::{
35 self,
36 TreeFile,
37};
38use anvil_job::{
39 ArtifactSpec,
40 CollectedArtifact,
41 JobSpec,
42 META_DIR,
43 META_TAR_CAP,
44 META_VALUE_CAP,
45 Sandbox,
46 Stored,
47 WORKDIR,
48};
49use bollard::{
50 Docker,
51 container::{
52 Config,
53 CreateContainerOptions,
54 DownloadFromContainerOptions,
55 LogsOptions,
56 RemoveContainerOptions,
57 StartContainerOptions,
58 UploadToContainerOptions,
59 WaitContainerOptions,
60 },
61 models::HostConfig,
62};
63use futures_util::StreamExt;
64use tokio::sync::mpsc::UnboundedReceiver;
65
66/// Turn a parsed pipeline into the job a runner receives.
67///
68/// The server's half of the split (see `docs/remote-runners.md`): image
69/// resolution, the allowlist check and script assembly all happen here. A
70/// runner therefore never parses `.anvil/ci.yml`, and cannot widen what it was
71/// permitted to run by reinterpreting one.
72fn build_job(
73 run_id: i64,
74 pipeline: &Pipeline,
75 cfg: &CiConfig,
76 env: &[(String, String)],
77) -> Result<JobSpec, String> {
78 // What the pipeline asked for, or the shared runner image when it omitted
79 // `image:` entirely.
80 let image = cfg.resolve_image(&pipeline.image);
81 if !cfg.image_allowed(image) {
82 return Err(format!(
83 "image {image} is not permitted by ci.allowed_images"
84 ));
85 }
86 Ok(JobSpec {
87 run_id,
88 image: image.to_string(),
89 // No source for this yet: the pipeline schema has no `platform:` key,
90 // so every job takes the runner daemon's native architecture, exactly
91 // as before. Honoured end to end the moment one is set.
92 platform: None,
93 script: build_script(pipeline),
94 env: env.to_vec(),
95 artifacts: pipeline
96 .artifacts
97 .iter()
98 .map(|a| ArtifactSpec {
99 name: a.name.clone(),
100 path: a.path.clone(),
101 browse: a.browse,
102 // The extractor commands themselves are already in the script;
103 // the runner only needs to know whether to look for output.
104 has_meta: !a.meta.is_empty(),
105 })
106 .collect(),
107 sandbox: Sandbox {
108 memory_mb: cfg.memory_mb,
109 cpus: cfg.cpus,
110 pids_limit: cfg.pids_limit,
111 timeout_secs: cfg.timeout_secs,
112 network: cfg.network,
113 run_as: cfg.run_as.clone(),
114 artifact_max_mb: cfg.artifact_max_mb,
115 artifact_run_max_mb: cfg.artifact_run_max_mb,
116 },
117 })
118}
119
120/// Assemble the single `sh -c` program a job runs.
121///
122/// The steps run in a subshell so the meta-extractor trailer still runs (and
123/// the original exit code is preserved) when a step fails — failure artifacts
124/// like test reports are the ones that matter most.
125fn build_script(pipeline: &Pipeline) -> String {
126 let mut script = String::from("(\nset -e\n");
127 for step in &pipeline.steps {
128 script.push_str("printf '\\n=== %s ===\\n' ");
129 script.push_str(&single_quote(step.label()));
130 script.push('\n');
131 script.push_str(&step.run);
132 script.push('\n');
133 }
134 script.push_str(")\nanvil_rc=$?\n");
135 for a in &pipeline.artifacts {
136 if a.meta.is_empty() {
137 continue;
138 }
139 // Names and keys are parse-time validated to [A-Za-z0-9._-]+, so they
140 // interpolate into the script safely.
141 script.push_str(&format!("mkdir -p {META_DIR}/{}\n", a.name));
142 for (key, cmd) in &a.meta {
143 script.push_str(&format!(
144 "{{\n{cmd}\n}} > {META_DIR}/{}/{key} 2>/dev/null || :\n",
145 a.name
146 ));
147 }
148 }
149 script.push_str("exit $anvil_rc\n");
150 script
151}
152
153/// Run the CI worker loop: recover interrupted runs, drain the queue, then
154/// process run ids as they arrive on `rx`. Runs one job at a time.
155pub async fn run_worker(app: App, mut rx: UnboundedReceiver<i64>) {
156 match ci::requeue_interrupted(&app.db).await {
157 Ok(ids) if !ids.is_empty() => {
158 tracing::info!("ci: requeued {} interrupted run(s)", ids.len())
159 }
160 Ok(_) => {}
161 Err(e) => tracing::error!("ci: requeue failed: {e}"),
162 }
163 match ci::queued_ids(&app.db).await {
164 Ok(ids) => {
165 for id in ids {
166 run_one(&app, id).await;
167 }
168 }
169 Err(e) => tracing::error!("ci: listing queued runs failed: {e}"),
170 }
171 tracing::info!("ci runner ready");
172 while let Some(id) = rx.recv().await {
173 run_one(&app, id).await;
174 }
175}
176
177async fn run_one(app: &App, run_id: i64) {
178 tracing::info!("ci: run {run_id} starting");
179 if let Err(e) = process(app, run_id).await {
180 tracing::error!("ci: run {run_id} errored: {e}");
181 let _ = ci::append_log(&app.db, run_id, &format!("\n[runner error] {e}\n")).await;
182 let _ = ci::finish(&app.db, run_id, ci::status::ERROR).await;
183 }
184}
185
186async fn process(app: &App, run_id: i64) -> Result<(), String> {
187 let run = ci::get(&app.db, run_id)
188 .await
189 .map_err(|e| e.to_string())?
190 .ok_or("run not found")?;
191 let repo = repos::find_by_id(&app.db, run.repo_id)
192 .await
193 .map_err(|e| e.to_string())?
194 .ok_or("repository not found")?;
195 let owner = users::find_by_id(&app.db, repo.owner_id)
196 .await
197 .map_err(|e| e.to_string())?
198 .ok_or("owner not found")?;
199 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
200
201 let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH)
202 .map_err(|e| e.to_string())?
203 .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?;
204 let pipeline =
205 ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?;
206
207 let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?;
208 let tar = build_tar(&files);
209
210 ci::mark_running(&app.db, run_id).await.ok();
211
212 let short = &run.commit[..run.commit.len().min(12)];
213 let mut log = format!(
214 "anvil ci · {}/{} · {} @ {short}\nimage: {}\n",
215 owner.username,
216 repo.name,
217 run.ref_name,
218 app.config.ci.resolve_image(&pipeline.image)
219 );
220
221 // Artifacts are collected into a scratch directory next to their final
222 // home (same filesystem, so the swap below is a rename), then moved into
223 // place only after the run finishes.
224 let artifacts_root = app.config.artifacts_dir();
225 let scratch = artifacts_root
226 .join(run.repo_id.to_string())
227 .join(format!(".collecting-{run_id}"));
228
229 // Secrets the pipeline asked for, from the in-memory vault. anvil holds no
230 // key that opens the stored envelopes, so an unlock must have happened
231 // (`anvild secret unlock`) or the run cannot proceed.
232 let env = match app.vault.take(run.repo_id, &pipeline.secrets) {
233 Ok(values) => values,
234 Err(missing) => {
235 log.push_str(&format!(
236 "\n[secrets unavailable: {}]\n\
237 This repository is sealed or was unlocked without them. Run:\n \
238 anvild secret unlock {}/{}\n",
239 missing.join(", "),
240 owner.username,
241 repo.name,
242 ));
243 ci::append_log(&app.db, run_id, &log).await.ok();
244 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
245 tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name);
246 return Ok(());
247 }
248 };
249 if !env.is_empty() {
250 log.push_str(&format!(
251 "secrets: {}\n",
252 env.iter()
253 .map(|(name, _)| name.as_str())
254 .collect::<Vec<_>>()
255 .join(", ")
256 ));
257 }
258
259 let mut sink = ScratchSink { scratch: &scratch };
260 let outcome = match build_job(run_id, &pipeline, &app.config.ci, &env) {
261 Ok(job) => execute(&job, tar, &mut log, &mut sink).await,
262 Err(e) => Err(e),
263 };
264 let (status, collected) = match outcome {
265 Ok((0, collected)) => (ci::status::SUCCESS, collected),
266 Ok((code, collected)) => {
267 log.push_str(&format!("\n[exited with status {code}]\n"));
268 (ci::status::FAILURE, collected)
269 }
270 Err(e) => {
271 log.push_str(&format!("\n[runner error] {e}\n"));
272 (ci::status::ERROR, Vec::new())
273 }
274 };
275
276 // Swap the collected set into place, replacing any earlier run's
277 // artifacts for this commit, then record the rows.
278 if !collected.is_empty() {
279 let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit);
280 let swap = async {
281 ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit)
282 .await
283 .map_err(|e| e.to_string())?;
284 if final_dir.exists() {
285 std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?;
286 }
287 std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?;
288 for c in &collected {
289 ci::add_artifact(
290 &app.db,
291 run_id,
292 run.repo_id,
293 &run.commit,
294 &c.name,
295 c.size,
296 c.is_dir,
297 c.browse,
298 &c.meta,
299 )
300 .await
301 .map_err(|e| e.to_string())?;
302 }
303 Ok::<_, String>(())
304 };
305 match swap.await {
306 Ok(()) => {
307 log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len()));
308 gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await;
309 }
310 Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")),
311 }
312 }
313 let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away
314
315 // A step that echoes its environment (`set -x`, `curl -v`, a failing
316 // command that prints its arguments) would otherwise publish the value on
317 // a page anyone with read access can see.
318 mask_secrets(&mut log, &env);
319 ci::append_log(&app.db, run_id, &log).await.ok();
320 ci::finish(&app.db, run_id, status).await.ok();
321 tracing::info!("ci: run {run_id} {status}");
322
323 // Continuous deployment: on a green run of the configured deploy repo's
324 // deploy branch, fire the redeploy webhook. Scoped to one repo by config —
325 // no other repository can trigger it, even with passing CI.
326 if status == ci::status::SUCCESS
327 && app
328 .config
329 .ci
330 .is_deploy_target(&owner.username, &repo.name, &run.ref_name)
331 {
332 deploy(app, &owner.username, &repo.name, &run).await;
333 }
334 Ok(())
335}
336
337/// Enforce `[ci] artifact_quota_mb` for one repository: while over budget,
338/// delete the oldest commit's artifacts (rows + directory). Branch-tip
339/// commits and the just-stored commit are pinned. Deterministic — runs after
340/// every artifact-producing run, no background sweeper. Best-effort: failures
341/// are logged, never failing the run.
342async fn gc_artifacts(
343 app: &App,
344 repo_id: i64,
345 repo_path: &Path,
346 keep_commit: &str,
347 log: &mut String,
348) {
349 let quota = match app.config.ci.artifact_quota_mb {
350 0 => return,
351 mb => mb_cap(mb),
352 };
353 let rows = match ci::artifacts_for_repo(&app.db, repo_id).await {
354 Ok(rows) => rows,
355 Err(e) => {
356 log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n"));
357 return;
358 }
359 };
360
361 // Per-commit totals and ages.
362 let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes)
363 let mut total: u64 = 0;
364 for row in &rows {
365 let entry = commits
366 .entry(row.commit.clone())
367 .or_insert((row.created_at, 0));
368 entry.0 = entry.0.min(row.created_at);
369 entry.1 += row.size;
370 total = total.saturating_add(row.size.max(0) as u64);
371 }
372 if total <= quota {
373 return;
374 }
375
376 let pinned = browse::branch_tips(repo_path).unwrap_or_default();
377 let mut victims: Vec<(i64, String, i64)> = commits
378 .into_iter()
379 .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit))
380 .map(|(commit, (oldest, bytes))| (oldest, commit, bytes))
381 .collect();
382 victims.sort();
383
384 for (_, commit, bytes) in victims {
385 if total <= quota {
386 break;
387 }
388 let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit);
389 if let Err(e) = std::fs::remove_dir_all(&dir) {
390 log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n"));
391 continue; // keep the rows; retried next run
392 }
393 if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await {
394 log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n"));
395 continue;
396 }
397 total = total.saturating_sub(bytes.max(0) as u64);
398 log.push_str(&format!(
399 "\n[artifact gc: dropped {} for old commit {}]\n",
400 fmt_mb(bytes),
401 &commit[..commit.len().min(12)]
402 ));
403 }
404}
405
406/// Bytes as a short MiB string for log lines.
407fn fmt_mb(bytes: i64) -> String {
408 format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0))
409}
410
411/// POST the configured deploy webhook. Best-effort: logs success/failure but
412/// never fails the run (CI already passed).
413async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) {
414 let cfg = &app.config.ci;
415 let body = serde_json::json!({
416 "repo": format!("{owner}/{name}"),
417 "ref": run.ref_name,
418 "commit": run.commit,
419 "run_id": run.id,
420 });
421 let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body);
422 if !cfg.deploy_secret.is_empty() {
423 req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret);
424 }
425 match req.send().await {
426 Ok(resp) if resp.status().is_success() => {
427 tracing::info!(
428 "ci: deploy webhook for {owner}/{name} accepted ({})",
429 resp.status()
430 )
431 }
432 Ok(resp) => tracing::error!(
433 "ci: deploy webhook for {owner}/{name} returned {}",
434 resp.status()
435 ),
436 Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"),
437 }
438}
439
440/// Where a collected artifact's bytes go.
441///
442/// Exists so the download loop can hand each tar off and drop it rather than
443/// buffering every artifact — `ci.artifact_run_max_mb` defaults to 512, which
444/// is not an amount to hold in RAM on the host anvil runs on. The in-process
445/// implementation writes to the scratch directory; a remote runner's uploads
446/// it. Returning [`Stored`] rather than `()` is what lets the caller charge
447/// the run budget the size that actually landed.
448#[async_trait::async_trait]
449pub trait ArtifactSink: Send {
450 async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String>;
451}
452
453/// The in-process sink: straight into the run's scratch directory, which
454/// `process` then renames into the commit's artifact directory.
455struct ScratchSink<'a> {
456 scratch: &'a Path,
457}
458
459#[async_trait::async_trait]
460impl ArtifactSink for ScratchSink<'_> {
461 async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String> {
462 // Created lazily rather than up front: a run whose artifacts all fail
463 // to download should leave no empty scratch directory behind.
464 std::fs::create_dir_all(self.scratch).map_err(|e| format!("creating scratch dir: {e}"))?;
465 let (size, is_dir) = store_artifact(spec, tar, self.scratch)?;
466 Ok(Stored { size, is_dir })
467 }
468}
469
470/// Replace every secret value in `log` with `***`.
471///
472/// Only values worth hiding: very short ones (a one-character secret) would
473/// mask half the log for no benefit, and are not credentials in practice.
474fn mask_secrets(log: &mut String, env: &[(String, String)]) {
475 for (_, value) in env {
476 if value.len() >= 4 && log.contains(value.as_str()) {
477 *log = log.replace(value.as_str(), "***");
478 }
479 }
480}
481
482/// Execute a job in a sandboxed container, streaming output into `log`.
483/// Returns the exit code and whatever artifacts `sink` accepted (none on
484/// timeout — the container is already gone).
485///
486/// The job container never sees the Docker socket and gets no mounts of any
487/// kind (the checkout is *uploaded*, not bind-mounted; artifacts are
488/// *downloaded* out the same way). All capabilities are dropped and
489/// `no-new-privileges` is set unconditionally; pids/memory/cpu caps, the
490/// wall-clock timeout, network access and the container user come from
491/// `spec.sandbox`. The image allowlist was applied when the spec was built —
492/// a runner receiving a spec does not get to widen it.
493pub async fn execute(
494 spec: &JobSpec,
495 tar: Vec<u8>,
496 log: &mut String,
497 sink: &mut dyn ArtifactSink,
498) -> Result<(i64, Vec<CollectedArtifact>), String> {
499 let platform = spec.platform.clone().unwrap_or_default();
500 let docker = anvil_docker::connect()?;
501 anvil_docker::ensure_image(&docker, &spec.image, &platform).await?;
502 let sb = &spec.sandbox;
503
504 // The sandbox. Limits of 0 mean "unlimited" and omit the corresponding cap.
505 let host_config = HostConfig {
506 cap_drop: Some(vec!["ALL".to_string()]),
507 security_opt: Some(vec!["no-new-privileges:true".to_string()]),
508 pids_limit: (sb.pids_limit > 0).then_some(sb.pids_limit),
509 memory: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024),
510 memory_swap: (sb.memory_mb > 0).then(|| sb.memory_mb * 1024 * 1024),
511 nano_cpus: (sb.cpus > 0.0).then_some((sb.cpus * 1e9) as i64),
512 network_mode: (!sb.network).then(|| "none".to_string()),
513 ..Default::default()
514 };
515 let config = Config {
516 image: Some(spec.image.clone()),
517 cmd: Some(vec![
518 "sh".to_string(),
519 "-c".to_string(),
520 spec.script.clone(),
521 ]),
522 env: (!spec.env.is_empty())
523 .then(|| spec.env.iter().map(|(k, v)| format!("{k}={v}")).collect()),
524 working_dir: Some(WORKDIR.to_string()),
525 user: (!sb.run_as.is_empty()).then(|| sb.run_as.clone()),
526 host_config: Some(host_config),
527 ..Default::default()
528 };
529 // An explicit platform needs the options struct; without one, pass None so
530 // the daemon picks its native architecture exactly as before.
531 let opts = spec.platform.as_ref().map(|p| CreateContainerOptions {
532 name: String::new(),
533 platform: Some(p.clone()),
534 });
535 let created = docker
536 .create_container(opts, config)
537 .await
538 .map_err(|e| format!("create container: {e}"))?;
539 let id = created.id;
540
541 // Upload the checkout (tar entries are under `workspace/`, extracted at `/`).
542 docker
543 .upload_to_container(
544 &id,
545 Some(UploadToContainerOptions {
546 path: "/".to_string(),
547 ..Default::default()
548 }),
549 tar.into(),
550 )
551 .await
552 .map_err(|e| format!("upload checkout: {e}"))?;
553
554 docker
555 .start_container(&id, None::<StartContainerOptions<String>>)
556 .await
557 .map_err(|e| format!("start container: {e}"))?;
558
559 // Stream logs and wait for the exit code, bounded by the wall-clock
560 // timeout. The container is force-removed on every path (which also kills
561 // a still-running job after a timeout).
562 let run = async {
563 let mut logs = docker.logs(
564 &id,
565 Some(LogsOptions::<String> {
566 follow: true,
567 stdout: true,
568 stderr: true,
569 ..Default::default()
570 }),
571 );
572 while let Some(item) = logs.next().await {
573 match item {
574 Ok(output) => log.push_str(&String::from_utf8_lossy(&output.into_bytes())),
575 Err(e) => {
576 log.push_str(&format!("\n[log stream error] {e}\n"));
577 break;
578 }
579 }
580 }
581
582 // Non-zero exit codes surface as a wait error in bollard.
583 let mut code = 0i64;
584 let mut wait = docker.wait_container(&id, None::<WaitContainerOptions<String>>);
585 while let Some(item) = wait.next().await {
586 match item {
587 Ok(resp) => code = resp.status_code,
588 Err(bollard::errors::Error::DockerContainerWaitError { code: c, .. }) => code = c,
589 Err(e) => return Err(format!("wait: {e}")),
590 }
591 }
592 Ok(code)
593 };
594 let result = match sb.timeout_secs {
595 0 => run.await,
596 secs => tokio::time::timeout(std::time::Duration::from_secs(secs), run)
597 .await
598 .unwrap_or_else(|_| Err(format!("job exceeded ci.timeout_secs ({secs}s); killed"))),
599 };
600
601 // Artifacts come out of the (now stopped) container before it is removed.
602 let collected = match &result {
603 Ok(_) if !spec.artifacts.is_empty() => {
604 collect_artifacts(&docker, &id, spec, log, sink).await
605 }
606 _ => Vec::new(),
607 };
608
609 let _ = docker
610 .remove_container(
611 &id,
612 Some(RemoveContainerOptions {
613 force: true,
614 ..Default::default()
615 }),
616 )
617 .await;
618
619 result.map(|code| (code, collected))
620}
621
622/// Collect the job's declared artifacts from the stopped container, handing
623/// each tar to `sink` as it is downloaded. Failures are per-artifact: each is
624/// logged and skipped, never failing the run.
625///
626/// One artifact is in memory at a time by construction — the tar is dropped
627/// once the sink has taken it — which is what keeps `artifact_run_max_mb`
628/// (512 MiB by default) a disk budget rather than a memory one.
629async fn collect_artifacts(
630 docker: &Docker,
631 id: &str,
632 job: &JobSpec,
633 log: &mut String,
634 sink: &mut dyn ArtifactSink,
635) -> Vec<CollectedArtifact> {
636 // Extractor outputs first: artifact name → key → value.
637 let mut metas: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new();
638 if job.artifacts.iter().any(|a| a.has_meta) {
639 match download_tar(docker, id, META_DIR, META_TAR_CAP).await {
640 Ok(Some(bytes)) => metas = parse_meta_tar(&bytes),
641 Ok(None) => log.push_str("\n[artifacts: extractor output exceeded its cap]\n"),
642 Err(e) => log.push_str(&format!("\n[artifacts: reading extractor output: {e}]\n")),
643 }
644 }
645
646 let per_artifact_cap = mb_cap(job.sandbox.artifact_max_mb);
647 let mut run_budget = mb_cap(job.sandbox.artifact_run_max_mb);
648 let mut collected = Vec::new();
649 for spec in &job.artifacts {
650 let cap = per_artifact_cap.min(run_budget);
651 let note = |log: &mut String, what: &str| {
652 log.push_str(&format!("\n[artifact {}: {what}]\n", spec.name));
653 };
654 let bytes = match download_tar(docker, id, &format!("{WORKDIR}/{}", spec.path), cap).await {
655 Ok(Some(bytes)) => bytes,
656 Ok(None) => {
657 note(log, "exceeds the size cap; skipped");
658 continue;
659 }
660 Err(e) => {
661 note(log, &format!("download failed ({e}); skipped"));
662 continue;
663 }
664 };
665 match sink.put(spec, &bytes).await {
666 Ok(Stored { size, is_dir }) => {
667 run_budget = run_budget.saturating_sub(size as u64);
668 let meta = metas.get(&spec.name).cloned().unwrap_or_default();
669 collected.push(CollectedArtifact {
670 name: spec.name.clone(),
671 size,
672 is_dir,
673 browse: spec.browse,
674 meta: serde_json::to_string(&meta).unwrap_or_else(|_| "{}".into()),
675 });
676 }
677 Err(e) => note(log, &format!("storing failed ({e}); skipped")),
678 }
679 }
680 collected
681}
682
683/// `0` (unlimited) → `u64::MAX`, otherwise MiB → bytes.
684fn mb_cap(mb: i64) -> u64 {
685 if mb <= 0 {
686 u64::MAX
687 } else {
688 mb as u64 * 1024 * 1024
689 }
690}
691
692/// Download `path` from the container as a tar, buffering at most `cap` bytes
693/// (`Ok(None)` when exceeded).
694async fn download_tar(
695 docker: &Docker,
696 id: &str,
697 path: &str,
698 cap: u64,
699) -> Result<Option<Vec<u8>>, String> {
700 let mut stream = docker.download_from_container(
701 id,
702 Some(DownloadFromContainerOptions {
703 path: path.to_string(),
704 }),
705 );
706 let mut buf = Vec::new();
707 while let Some(chunk) = stream.next().await {
708 let chunk = chunk.map_err(|e| e.to_string())?;
709 if (buf.len() + chunk.len()) as u64 > cap {
710 return Ok(None);
711 }
712 buf.extend_from_slice(&chunk);
713 }
714 Ok(Some(buf))
715}
716
717/// Parse the extractor-output tar (`anvil-meta/<artifact>/<key>` files) into
718/// artifact → key → trimmed value.
719fn parse_meta_tar(bytes: &[u8]) -> BTreeMap<String, BTreeMap<String, String>> {
720 let mut out: BTreeMap<String, BTreeMap<String, String>> = BTreeMap::new();
721 let mut archive = tar::Archive::new(bytes);
722 let Ok(entries) = archive.entries() else {
723 return out;
724 };
725 for entry in entries.flatten() {
726 if !entry.header().entry_type().is_file() {
727 continue;
728 }
729 let Ok(path) = entry.path() else { continue };
730 // anvil-meta/<artifact>/<key>
731 let parts: Vec<String> = path
732 .components()
733 .skip(1)
734 .map(|c| c.as_os_str().to_string_lossy().into_owned())
735 .collect();
736 let [artifact, key] = parts.as_slice() else {
737 continue;
738 };
739 let (artifact, key) = (artifact.clone(), key.clone());
740 let mut value = String::new();
741 let _ = entry.take(META_VALUE_CAP as u64).read_to_string(&mut value);
742 let value = value.trim().to_string();
743 if !value.is_empty() {
744 out.entry(artifact).or_default().insert(key, value);
745 }
746 }
747 out
748}
749
750/// Write one downloaded artifact tar into `scratch`, returning (size, is_dir).
751///
752/// - a file is stored as-is at `scratch/<name>`
753/// - a directory with `browse` is extracted under `scratch/<name>/`
754/// - any other directory is stored compressed at `scratch/<name>.tar.gz`
755fn store_artifact(
756 spec: &ArtifactSpec,
757 tar_bytes: &[u8],
758 scratch: &Path,
759) -> Result<(i64, bool), String> {
760 // The docker archive endpoint roots entries at the requested item's
761 // basename; its first entry tells file from directory.
762 let mut archive = tar::Archive::new(tar_bytes);
763 let mut entries = archive.entries().map_err(|e| e.to_string())?;
764 let first = entries
765 .next()
766 .ok_or("empty archive")?
767 .map_err(|e| e.to_string())?;
768 let is_dir = first.header().entry_type().is_dir();
769
770 if !is_dir {
771 let mut entry = first;
772 let mut content = Vec::new();
773 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
774 std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?;
775 return Ok((content.len() as i64, false));
776 }
777
778 if !spec.browse {
779 let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name)))
780 .map_err(|e| e.to_string())?;
781 let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default());
782 std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?;
783 let file = enc.finish().map_err(|e| e.to_string())?;
784 let size = file.metadata().map_err(|e| e.to_string())?.len();
785 return Ok((size as i64, true));
786 }
787
788 // Browsable: extract regular files under scratch/<name>/, stripping the
789 // basename prefix. Entry paths come from docker's tar of a real
790 // filesystem, but stay defensive: relative components only, no links.
791 let root = scratch.join(&spec.name);
792 let mut total = 0i64;
793 let mut archive = tar::Archive::new(tar_bytes);
794 for entry in archive.entries().map_err(|e| e.to_string())?.flatten() {
795 if !entry.header().entry_type().is_file() {
796 continue;
797 }
798 let Ok(path) = entry.path() else { continue };
799 let mut rel = std::path::PathBuf::new();
800 let mut ok = true;
801 for c in path.components().skip(1) {
802 match c {
803 std::path::Component::Normal(p) => rel.push(p),
804 _ => {
805 ok = false;
806 break;
807 }
808 }
809 }
810 if !ok || rel.as_os_str().is_empty() {
811 continue;
812 }
813 let dest = root.join(&rel);
814 if let Some(parent) = dest.parent() {
815 std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
816 }
817 let mut entry = entry;
818 let mut content = Vec::new();
819 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
820 std::fs::write(&dest, &content).map_err(|e| e.to_string())?;
821 total += content.len() as i64;
822 }
823 let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists
824 Ok((total, true))
825}
826
827/// Build an uncompressed tar of the checkout, rooted at `workspace/` so it
828/// extracts to `/workspace` when uploaded to the container root.
829fn build_tar(files: &[TreeFile]) -> Vec<u8> {
830 let mut builder = tar::Builder::new(Vec::new());
831 for f in files {
832 let mut header = tar::Header::new_gnu();
833 header.set_size(f.content.len() as u64);
834 header.set_mode(if f.executable { 0o755 } else { 0o644 });
835 // append_data sets the path and checksum.
836 let _ = builder.append_data(
837 &mut header,
838 format!("workspace/{}", f.path),
839 f.content.as_slice(),
840 );
841 }
842 builder.into_inner().unwrap_or_default()
843}
844
845/// Single-quote a string for safe interpolation into a shell command.
846fn single_quote(s: &str) -> String {
847 format!("'{}'", s.replace('\'', "'\\''"))
848}
849
850#[cfg(test)]
851mod tests {
852 use super::*;
853
854 /// Build a tar the way docker's archive endpoint does: entries rooted at
855 /// the requested item's basename.
856 fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> {
857 let mut b = tar::Builder::new(Vec::new());
858 for (path, content) in entries {
859 let mut h = tar::Header::new_gnu();
860 match content {
861 Some(c) => {
862 h.set_size(c.len() as u64);
863 h.set_mode(0o644);
864 h.set_entry_type(tar::EntryType::Regular);
865 b.append_data(&mut h, path, c.as_bytes()).unwrap();
866 }
867 None => {
868 h.set_size(0);
869 h.set_mode(0o755);
870 h.set_entry_type(tar::EntryType::Directory);
871 b.append_data(&mut h, path, std::io::empty()).unwrap();
872 }
873 }
874 }
875 b.into_inner().unwrap()
876 }
877
878 fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec {
879 ArtifactSpec {
880 name: name.into(),
881 path: path.into(),
882 browse,
883 has_meta: false,
884 }
885 }
886
887 #[test]
888 fn stores_a_file_artifact_as_is() {
889 let dir = tempfile::tempdir().unwrap();
890 let tar = tar_of(&[("anvild", Some("ELF..."))]);
891 let (size, is_dir) = store_artifact(
892 &spec("bin", "target/release/anvild", false),
893 &tar,
894 dir.path(),
895 )
896 .unwrap();
897 assert!(!is_dir);
898 assert_eq!(size, 6);
899 assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF...");
900 }
901
902 #[test]
903 fn stores_a_directory_artifact_as_tar_gz() {
904 let dir = tempfile::tempdir().unwrap();
905 let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]);
906 let (size, is_dir) =
907 store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap();
908 assert!(is_dir);
909 let stored = dir.path().join("doc.tar.gz");
910 assert_eq!(size, stored.metadata().unwrap().len() as i64);
911 // Round-trips through gzip back to the original tar bytes.
912 let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap());
913 let mut bytes = Vec::new();
914 gz.read_to_end(&mut bytes).unwrap();
915 assert_eq!(bytes, tar);
916 }
917
918 #[test]
919 fn extracts_a_browsable_directory_artifact() {
920 let dir = tempfile::tempdir().unwrap();
921 let tar = tar_of(&[
922 ("doc", None),
923 ("doc/index.html", Some("<html>")),
924 ("doc/sub", None),
925 ("doc/sub/page.html", Some("<p>")),
926 ]);
927 let (size, is_dir) =
928 store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap();
929 assert!(is_dir);
930 assert_eq!(size, 6 + 3);
931 let root = dir.path().join("doc");
932 assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>");
933 assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>");
934 }
935
936 #[test]
937 fn parses_meta_tar_with_trimmed_capped_values() {
938 let tar = tar_of(&[
939 ("anvil-meta", None),
940 ("anvil-meta/bin", None),
941 ("anvil-meta/bin/version", Some("anvild 0.0.0\n")),
942 ("anvil-meta/bin/empty", Some(" \n")),
943 ("anvil-meta/doc", None),
944 ("anvil-meta/doc/pages", Some("42")),
945 ]);
946 let metas = parse_meta_tar(&tar);
947 assert_eq!(metas["bin"]["version"], "anvild 0.0.0");
948 assert_eq!(metas["doc"]["pages"], "42");
949 assert!(!metas["bin"].contains_key("empty"), "blank values dropped");
950 }
951}