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