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