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