anvilsign in

collin/anvil

1//! CI dispatch: decides what a queued [`anvil_core::ci`] run should do, hands
2//! it to a runner that dialled in, and records what came back.
3//!
4//! anvil does **not** execute jobs. It used to — this crate was the runner —
5//! but a release build needs more RAM than the host it runs on has, so
6//! execution moved to `anvil-worker` on a machine that can afford it. See
7//! `docs/remote-runners.md`.
8//!
9//! What stays here is everything a runner must not be trusted with: resolving
10//! the repo, materializing the commit's tree, parsing `.anvil/ci.yml`, checking
11//! the image allowlist, opening the secret vault, assembling the shell script,
12//! and deciding where artifacts land on disk. A runner receives an image, a
13//! script, a tar and some limits ([`anvil_job::JobSpec`]) and can neither widen
14//! its sandbox nor reinterpret a pipeline.
15
16use std::{
17 collections::BTreeMap,
18 io::Read,
19 path::Path,
20};
21
22use anvil_core::{
23 App,
24 ci::{
25 self,
26 Pipeline,
27 },
28 config::CiConfig,
29 repos,
30 storage,
31 users,
32};
33use anvil_git::browse::{
34 self,
35 TreeFile,
36};
37use anvil_job::{
38 ArtifactSpec,
39 JobResult,
40 JobSpec,
41 META_DIR,
42 Sandbox,
43 Stored,
44 mb_cap,
45};
46use tokio::sync::mpsc::UnboundedReceiver;
47
48/// Turn a parsed pipeline into the job a runner receives.
49///
50/// The server's half of the split (see `docs/remote-runners.md`): image
51/// resolution, the allowlist check and script assembly all happen here. A
52/// runner therefore never parses `.anvil/ci.yml`, and cannot widen what it was
53/// permitted to run by reinterpreting one.
54fn build_job(
55 run_id: i64,
56 pipeline: &Pipeline,
57 cfg: &CiConfig,
58 env: &[(String, String)],
59) -> Result<JobSpec, String> {
60 // What the pipeline asked for, or the shared runner image when it omitted
61 // `image:` entirely.
62 let image = cfg.resolve_image(&pipeline.image);
63 if !cfg.image_allowed(image) {
64 return Err(format!(
65 "image {image} is not permitted by ci.allowed_images"
66 ));
67 }
68 Ok(JobSpec {
69 run_id,
70 image: image.to_string(),
71 // No source for this yet: the pipeline schema has no `platform:` key,
72 // so every job takes the runner daemon's native architecture, exactly
73 // as before. Honoured end to end the moment one is set.
74 platform: None,
75 script: build_script(pipeline),
76 env: env.to_vec(),
77 artifacts: pipeline
78 .artifacts
79 .iter()
80 .map(|a| ArtifactSpec {
81 name: a.name.clone(),
82 path: a.path.clone(),
83 browse: a.browse,
84 // The extractor commands themselves are already in the script;
85 // the runner only needs to know whether to look for output.
86 has_meta: !a.meta.is_empty(),
87 })
88 .collect(),
89 sandbox: Sandbox {
90 memory_mb: cfg.memory_mb,
91 cpus: cfg.cpus,
92 pids_limit: cfg.pids_limit,
93 timeout_secs: cfg.timeout_secs,
94 network: cfg.network,
95 run_as: cfg.run_as.clone(),
96 artifact_max_mb: cfg.artifact_max_mb,
97 artifact_run_max_mb: cfg.artifact_run_max_mb,
98 },
99 })
100}
101
102/// Assemble the single `sh -c` program a job runs.
103///
104/// The steps run in a subshell so the meta-extractor trailer still runs (and
105/// the original exit code is preserved) when a step fails — failure artifacts
106/// like test reports are the ones that matter most.
107fn build_script(pipeline: &Pipeline) -> String {
108 let mut script = String::from("(\nset -e\n");
109 for step in &pipeline.steps {
110 script.push_str("printf '\\n=== %s ===\\n' ");
111 script.push_str(&single_quote(step.label()));
112 script.push('\n');
113 script.push_str(&step.run);
114 script.push('\n');
115 }
116 script.push_str(")\nanvil_rc=$?\n");
117 for a in &pipeline.artifacts {
118 if a.meta.is_empty() {
119 continue;
120 }
121 // Names and keys are parse-time validated to [A-Za-z0-9._-]+, so they
122 // interpolate into the script safely.
123 script.push_str(&format!("mkdir -p {META_DIR}/{}\n", a.name));
124 for (key, cmd) in &a.meta {
125 script.push_str(&format!(
126 "{{\n{cmd}\n}} > {META_DIR}/{}/{key} 2>/dev/null || :\n",
127 a.name
128 ));
129 }
130 }
131 script.push_str("exit $anvil_rc\n");
132 script
133}
134
135/// Dispatch queued runs to runners: recover interrupted runs, then wake parked
136/// claim requests as work arrives, and requeue jobs whose runner went away.
137///
138/// anvild no longer executes anything (see `docs/remote-runners.md`). This
139/// loop used to *be* the runner; now it only decides that a run is ready to be
140/// claimed. `rx` still carries newly-enqueued run ids from every push path, so
141/// nothing that enqueues had to change.
142pub async fn run_dispatcher(app: App, mut rx: UnboundedReceiver<i64>) {
143 match ci::requeue_interrupted(&app.db).await {
144 Ok(ids) if !ids.is_empty() => {
145 tracing::info!("ci: requeued {} interrupted run(s)", ids.len())
146 }
147 Ok(_) => {}
148 Err(e) => tracing::error!("ci: requeue failed: {e}"),
149 }
150
151 tokio::spawn(sweep_expired_leases(app.clone()));
152
153 if app.config.ci.runner_token.is_empty() {
154 tracing::warn!(
155 "ci: [ci] runner_token is empty, so no runner can authenticate — \
156 queued runs will sit until one is set"
157 );
158 }
159 tracing::info!("ci dispatcher ready");
160
161 // Anything already queued is claimable the moment a runner asks; the wake
162 // is only for runners already parked in a long poll.
163 app.jobs.wake();
164 while rx.recv().await.is_some() {
165 app.jobs.wake();
166 }
167}
168
169/// Requeue runs whose runner stopped heartbeating.
170///
171/// The case anvild's own crash never had: with execution in-process, a job
172/// died exactly when the process did, and `requeue_interrupted` at startup
173/// covered it. A remote runner can vanish while anvil stays up.
174///
175/// Note this can double-run a job whose runner is alive but unreachable — the
176/// container keeps going, and the requeued run may be claimed elsewhere. CI
177/// steps are assumed idempotent, and one job at a time makes it unlikely, but
178/// it is the honest failure mode of a lease without fencing.
179async fn sweep_expired_leases(app: App) {
180 let mut tick = tokio::time::interval(anvil_core::jobs::LEASE_TTL / 4);
181 loop {
182 tick.tick().await;
183 for (run_id, runner) in app.jobs.take_expired() {
184 tracing::warn!("ci: run {run_id} lease expired (runner {runner}); requeueing");
185 if let Err(e) = ci::requeue(&app.db, run_id).await {
186 tracing::error!("ci: requeueing {run_id} failed: {e}");
187 continue;
188 }
189 app.jobs.wake();
190 }
191 }
192}
193
194/// Hand the oldest queued run to `runner`, or `None` when there is nothing to
195/// do.
196///
197/// Runs that cannot be dispatched at all — an unparseable pipeline, a sealed
198/// vault, an image the allowlist forbids — are finished here rather than
199/// handed out, and the next queued run is tried. That keeps a single bad
200/// pipeline from wedging the queue.
201pub async fn claim_next(app: &App, runner: &str) -> Option<JobSpec> {
202 let queued = match ci::queued_ids(&app.db).await {
203 Ok(ids) => ids,
204 Err(e) => {
205 tracing::error!("ci: listing queued runs failed: {e}");
206 return None;
207 }
208 };
209 for run_id in queued {
210 // Skip anything a concurrent claim already took: `queued_ids` reads
211 // committed state, and `prepare` marks running.
212 if app.jobs.holds_any(run_id) {
213 continue;
214 }
215 match prepare(app, run_id, runner).await {
216 Ok(Some(job)) => return Some(job),
217 Ok(None) => continue,
218 Err(e) => {
219 tracing::error!("ci: preparing run {run_id} failed: {e}");
220 fail_run(app, run_id, &format!("\n[dispatch error] {e}\n")).await;
221 }
222 }
223 }
224 None
225}
226
227/// Everything anvil does before a job leaves the building: resolve the repo,
228/// parse the pipeline, open the vault, build the spec, take the lease.
229///
230/// `Ok(None)` means the run was finished here and should not be dispatched.
231async fn prepare(app: &App, run_id: i64, runner: &str) -> Result<Option<JobSpec>, String> {
232 let run = ci::get(&app.db, run_id)
233 .await
234 .map_err(|e| e.to_string())?
235 .ok_or("run not found")?;
236 let repo = repos::find_by_id(&app.db, run.repo_id)
237 .await
238 .map_err(|e| e.to_string())?
239 .ok_or("repository not found")?;
240 let owner = users::find_by_id(&app.db, repo.owner_id)
241 .await
242 .map_err(|e| e.to_string())?
243 .ok_or("owner not found")?;
244 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
245
246 let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH)
247 .map_err(|e| e.to_string())?
248 .ok_or_else(|| format!("{} missing at {}", ci::PIPELINE_PATH, run.commit))?;
249 let pipeline =
250 ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?;
251
252 let short = &run.commit[..run.commit.len().min(12)];
253 let mut log = format!(
254 "anvil ci · {}/{} · {} @ {short}\nrunner: {runner}\nimage: {}\n",
255 owner.username,
256 repo.name,
257 run.ref_name,
258 app.config.ci.resolve_image(&pipeline.image)
259 );
260
261 // Secrets the pipeline asked for, from the in-memory vault. anvil holds no
262 // key that opens the stored envelopes, so an unlock must have happened
263 // (`anvild secret unlock`) or the run cannot proceed.
264 let env = match app.vault.take(run.repo_id, &pipeline.secrets) {
265 Ok(values) => values,
266 Err(missing) => {
267 log.push_str(&format!(
268 "\n[secrets unavailable: {}]\n\
269 This repository is sealed or was unlocked without them. Run:\n \
270 anvild secret unlock {}/{}\n",
271 missing.join(", "),
272 owner.username,
273 repo.name,
274 ));
275 ci::append_log(&app.db, run_id, &log).await.ok();
276 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
277 tracing::warn!("ci: run {run_id} needs secrets but {} is sealed", repo.name);
278 return Ok(None);
279 }
280 };
281 if !env.is_empty() {
282 log.push_str(&format!(
283 "secrets: {}\n",
284 env.iter()
285 .map(|(name, _)| name.as_str())
286 .collect::<Vec<_>>()
287 .join(", ")
288 ));
289 }
290
291 let job = match build_job(run_id, &pipeline, &app.config.ci, &env) {
292 Ok(job) => job,
293 Err(e) => {
294 log.push_str(&format!("\n[runner error] {e}\n"));
295 ci::append_log(&app.db, run_id, &log).await.ok();
296 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
297 return Ok(None);
298 }
299 };
300
301 // The header is written now rather than kept in memory until the result
302 // lands, so a run that is visibly `running` has something to show.
303 ci::append_log(&app.db, run_id, &log).await.ok();
304 ci::mark_running(&app.db, run_id).await.ok();
305 app.jobs.claim(run_id, runner, env);
306 tracing::info!("ci: run {run_id} claimed by {runner}");
307 Ok(Some(job))
308}
309
310/// The checkout a claimed job runs against, as an uncompressed tar.
311///
312/// Rebuilt from the commit on demand rather than stashed at claim time: it is
313/// a pure function of the commit, so a runner retrying the fetch is free, and
314/// anvil holds no per-job buffer on a host where memory is the scarce thing.
315pub async fn checkout_tar(app: &App, run_id: i64) -> Result<Vec<u8>, String> {
316 let (_, _, repo_path, run) = resolve(app, run_id).await?;
317 let files = browse::read_tree_files(&repo_path, &run.commit).map_err(|e| e.to_string())?;
318 Ok(build_tar(&files))
319}
320
321/// The declared spec for one artifact of a run.
322///
323/// Re-derived from the pipeline rather than trusted from the upload request:
324/// `browse` decides whether a tarball is extracted into a servable directory
325/// tree, which is not a choice a runner should get to make.
326pub async fn artifact_spec(
327 app: &App,
328 run_id: i64,
329 name: &str,
330) -> Result<Option<ArtifactSpec>, String> {
331 let (_, _, repo_path, run) = resolve(app, run_id).await?;
332 let yaml = browse::read_blob(&repo_path, &run.commit, ci::PIPELINE_PATH)
333 .map_err(|e| e.to_string())?
334 .ok_or("pipeline missing")?;
335 let pipeline =
336 ci::parse_pipeline(&String::from_utf8_lossy(&yaml)).map_err(|e| e.to_string())?;
337 Ok(pipeline
338 .artifacts
339 .iter()
340 .find(|a| a.name == name)
341 .map(|a| ArtifactSpec {
342 name: a.name.clone(),
343 path: a.path.clone(),
344 browse: a.browse,
345 has_meta: !a.meta.is_empty(),
346 }))
347}
348
349/// Store one uploaded artifact tar into the run's scratch directory.
350///
351/// Scratch sits next to the artifacts' final home (same filesystem, so
352/// [`finish_run`]'s swap is a rename) and is only moved into place once the
353/// run reports in.
354pub fn store_upload(
355 app: &App,
356 run: &anvil_core::CiRun,
357 spec: &ArtifactSpec,
358 tar: &[u8],
359) -> Result<Stored, String> {
360 let scratch = scratch_dir(app, run);
361 std::fs::create_dir_all(&scratch).map_err(|e| format!("creating scratch dir: {e}"))?;
362 let (size, is_dir) = store_artifact(spec, tar, &scratch)?;
363 Ok(Stored { size, is_dir })
364}
365
366/// Where a run's artifacts accumulate before the swap.
367fn scratch_dir(app: &App, run: &anvil_core::CiRun) -> std::path::PathBuf {
368 app.config
369 .artifacts_dir()
370 .join(run.repo_id.to_string())
371 .join(format!(".collecting-{}", run.id))
372}
373
374/// Record a finished job: swap its artifacts into place, mask and store the
375/// log, set the status, and fire the deploy webhook if it qualifies.
376pub async fn finish_run(app: &App, run_id: i64, result: JobResult) -> Result<(), String> {
377 let (owner, repo, repo_path, run) = resolve(app, run_id).await?;
378 let scratch = scratch_dir(app, &run);
379 let mut log = result.log;
380
381 let (status, collected) = match result.runner_error {
382 Some(e) => {
383 log.push_str(&format!("\n[runner error] {e}\n"));
384 (ci::status::ERROR, Vec::new())
385 }
386 None if result.exit_code == 0 => (ci::status::SUCCESS, result.artifacts),
387 None => {
388 log.push_str(&format!("\n[exited with status {}]\n", result.exit_code));
389 (ci::status::FAILURE, result.artifacts)
390 }
391 };
392
393 // Swap the collected set into place, replacing any earlier run's
394 // artifacts for this commit, then record the rows.
395 if !collected.is_empty() {
396 let artifacts_root = app.config.artifacts_dir();
397 let final_dir = storage::artifact_commit_dir(&artifacts_root, run.repo_id, &run.commit);
398 let swap = async {
399 ci::delete_artifacts_for_commit(&app.db, run.repo_id, &run.commit)
400 .await
401 .map_err(|e| e.to_string())?;
402 if final_dir.exists() {
403 std::fs::remove_dir_all(&final_dir).map_err(|e| e.to_string())?;
404 }
405 std::fs::rename(&scratch, &final_dir).map_err(|e| e.to_string())?;
406 for c in &collected {
407 ci::add_artifact(
408 &app.db,
409 run_id,
410 run.repo_id,
411 &run.commit,
412 &c.name,
413 c.size,
414 c.is_dir,
415 c.browse,
416 &c.meta,
417 )
418 .await
419 .map_err(|e| e.to_string())?;
420 }
421 Ok::<_, String>(())
422 };
423 match swap.await {
424 Ok(()) => {
425 log.push_str(&format!("\n[collected {} artifact(s)]\n", collected.len()));
426 gc_artifacts(app, run.repo_id, &repo_path, &run.commit, &mut log).await;
427 }
428 Err(e) => log.push_str(&format!("\n[storing artifacts failed: {e}]\n")),
429 }
430 }
431 let _ = std::fs::remove_dir_all(&scratch); // no-op when renamed away
432
433 // A step that echoes its environment (`set -x`, `curl -v`, a failing
434 // command that prints its arguments) would otherwise publish the value on
435 // a page anyone with read access can see. The values come from the lease
436 // rather than the vault, which may have re-sealed while the job ran.
437 mask_secrets(&mut log, &app.jobs.secrets_for(run_id));
438 ci::append_log(&app.db, run_id, &log).await.ok();
439 ci::finish(&app.db, run_id, status).await.ok();
440 app.jobs.release(run_id);
441 tracing::info!("ci: run {run_id} {status}");
442
443 // Continuous deployment: on a green run of the configured deploy repo's
444 // deploy branch, fire the redeploy webhook. Scoped to one repo by config —
445 // no other repository can trigger it, even with passing CI.
446 if status == ci::status::SUCCESS && app.config.ci.is_deploy_target(&owner, &repo, &run.ref_name)
447 {
448 deploy(app, &owner, &repo, &run).await;
449 }
450 Ok(())
451}
452
453/// Finish a run that never got far enough to produce a result.
454async fn fail_run(app: &App, run_id: i64, note: &str) {
455 ci::append_log(&app.db, run_id, note).await.ok();
456 ci::finish(&app.db, run_id, ci::status::ERROR).await.ok();
457 app.jobs.release(run_id);
458}
459
460/// (owner username, repo name, repo path, run) for a run id.
461async fn resolve(
462 app: &App,
463 run_id: i64,
464) -> Result<(String, String, std::path::PathBuf, anvil_core::CiRun), String> {
465 let run = ci::get(&app.db, run_id)
466 .await
467 .map_err(|e| e.to_string())?
468 .ok_or("run not found")?;
469 let repo = repos::find_by_id(&app.db, run.repo_id)
470 .await
471 .map_err(|e| e.to_string())?
472 .ok_or("repository not found")?;
473 let owner = users::find_by_id(&app.db, repo.owner_id)
474 .await
475 .map_err(|e| e.to_string())?
476 .ok_or("owner not found")?;
477 let repo_path = storage::repo_path(&app.config.repositories_dir(), &owner.username, &repo.name);
478 Ok((owner.username, repo.name, repo_path, run))
479}
480
481/// Enforce `[ci] artifact_quota_mb` for one repository: while over budget,
482/// delete the oldest commit's artifacts (rows + directory). Branch-tip
483/// commits and the just-stored commit are pinned. Deterministic — runs after
484/// every artifact-producing run, no background sweeper. Best-effort: failures
485/// are logged, never failing the run.
486async fn gc_artifacts(
487 app: &App,
488 repo_id: i64,
489 repo_path: &Path,
490 keep_commit: &str,
491 log: &mut String,
492) {
493 let quota = match app.config.ci.artifact_quota_mb {
494 0 => return,
495 mb => mb_cap(mb),
496 };
497 let rows = match ci::artifacts_for_repo(&app.db, repo_id).await {
498 Ok(rows) => rows,
499 Err(e) => {
500 log.push_str(&format!("\n[artifact gc: listing failed: {e}]\n"));
501 return;
502 }
503 };
504
505 // Per-commit totals and ages.
506 let mut commits: BTreeMap<String, (i64, i64)> = BTreeMap::new(); // commit → (oldest created_at, bytes)
507 let mut total: u64 = 0;
508 for row in &rows {
509 let entry = commits
510 .entry(row.commit.clone())
511 .or_insert((row.created_at, 0));
512 entry.0 = entry.0.min(row.created_at);
513 entry.1 += row.size;
514 total = total.saturating_add(row.size.max(0) as u64);
515 }
516 if total <= quota {
517 return;
518 }
519
520 let pinned = browse::branch_tips(repo_path).unwrap_or_default();
521 let mut victims: Vec<(i64, String, i64)> = commits
522 .into_iter()
523 .filter(|(commit, _)| commit != keep_commit && !pinned.contains(commit))
524 .map(|(commit, (oldest, bytes))| (oldest, commit, bytes))
525 .collect();
526 victims.sort();
527
528 for (_, commit, bytes) in victims {
529 if total <= quota {
530 break;
531 }
532 let dir = storage::artifact_commit_dir(&app.config.artifacts_dir(), repo_id, &commit);
533 if let Err(e) = std::fs::remove_dir_all(&dir) {
534 log.push_str(&format!("\n[artifact gc: removing {commit}: {e}]\n"));
535 continue; // keep the rows; retried next run
536 }
537 if let Err(e) = ci::delete_artifacts_for_commit(&app.db, repo_id, &commit).await {
538 log.push_str(&format!("\n[artifact gc: forgetting {commit}: {e}]\n"));
539 continue;
540 }
541 total = total.saturating_sub(bytes.max(0) as u64);
542 log.push_str(&format!(
543 "\n[artifact gc: dropped {} for old commit {}]\n",
544 fmt_mb(bytes),
545 &commit[..commit.len().min(12)]
546 ));
547 }
548}
549
550/// Bytes as a short MiB string for log lines.
551fn fmt_mb(bytes: i64) -> String {
552 format!("{:.1} MiB", bytes.max(0) as f64 / (1024.0 * 1024.0))
553}
554
555/// POST the configured deploy webhook. Best-effort: logs success/failure but
556/// never fails the run (CI already passed).
557async fn deploy(app: &App, owner: &str, name: &str, run: &anvil_core::CiRun) {
558 let cfg = &app.config.ci;
559 let body = serde_json::json!({
560 "repo": format!("{owner}/{name}"),
561 "ref": run.ref_name,
562 "commit": run.commit,
563 "run_id": run.id,
564 });
565 let mut req = reqwest::Client::new().post(&cfg.deploy_webhook).json(&body);
566 if !cfg.deploy_secret.is_empty() {
567 req = req.header("X-Anvil-Deploy-Secret", &cfg.deploy_secret);
568 }
569 match req.send().await {
570 Ok(resp) if resp.status().is_success() => {
571 tracing::info!(
572 "ci: deploy webhook for {owner}/{name} accepted ({})",
573 resp.status()
574 )
575 }
576 Ok(resp) => tracing::error!(
577 "ci: deploy webhook for {owner}/{name} returned {}",
578 resp.status()
579 ),
580 Err(e) => tracing::error!("ci: deploy webhook for {owner}/{name} failed: {e}"),
581 }
582}
583
584/// Replace every secret value in `log` with `***`.
585///
586/// Only values worth hiding: very short ones (a one-character secret) would
587/// mask half the log for no benefit, and are not credentials in practice.
588fn mask_secrets(log: &mut String, env: &[(String, String)]) {
589 for (_, value) in env {
590 if value.len() >= 4 && log.contains(value.as_str()) {
591 *log = log.replace(value.as_str(), "***");
592 }
593 }
594}
595
596/// Write one downloaded artifact tar into `scratch`, returning (size, is_dir).
597///
598/// - a file is stored as-is at `scratch/<name>`
599/// - a directory with `browse` is extracted under `scratch/<name>/`
600/// - any other directory is stored compressed at `scratch/<name>.tar.gz`
601fn store_artifact(
602 spec: &ArtifactSpec,
603 tar_bytes: &[u8],
604 scratch: &Path,
605) -> Result<(i64, bool), String> {
606 // The docker archive endpoint roots entries at the requested item's
607 // basename; its first entry tells file from directory.
608 let mut archive = tar::Archive::new(tar_bytes);
609 let mut entries = archive.entries().map_err(|e| e.to_string())?;
610 let first = entries
611 .next()
612 .ok_or("empty archive")?
613 .map_err(|e| e.to_string())?;
614 let is_dir = first.header().entry_type().is_dir();
615
616 if !is_dir {
617 let mut entry = first;
618 let mut content = Vec::new();
619 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
620 std::fs::write(scratch.join(&spec.name), &content).map_err(|e| e.to_string())?;
621 return Ok((content.len() as i64, false));
622 }
623
624 if !spec.browse {
625 let file = std::fs::File::create(scratch.join(format!("{}.tar.gz", spec.name)))
626 .map_err(|e| e.to_string())?;
627 let mut enc = flate2::write::GzEncoder::new(file, flate2::Compression::default());
628 std::io::Write::write_all(&mut enc, tar_bytes).map_err(|e| e.to_string())?;
629 let file = enc.finish().map_err(|e| e.to_string())?;
630 let size = file.metadata().map_err(|e| e.to_string())?.len();
631 return Ok((size as i64, true));
632 }
633
634 // Browsable: extract regular files under scratch/<name>/, stripping the
635 // basename prefix. Entry paths come from docker's tar of a real
636 // filesystem, but stay defensive: relative components only, no links.
637 let root = scratch.join(&spec.name);
638 let mut total = 0i64;
639 let mut archive = tar::Archive::new(tar_bytes);
640 for entry in archive.entries().map_err(|e| e.to_string())?.flatten() {
641 if !entry.header().entry_type().is_file() {
642 continue;
643 }
644 let Ok(path) = entry.path() else { continue };
645 let mut rel = std::path::PathBuf::new();
646 let mut ok = true;
647 for c in path.components().skip(1) {
648 match c {
649 std::path::Component::Normal(p) => rel.push(p),
650 _ => {
651 ok = false;
652 break;
653 }
654 }
655 }
656 if !ok || rel.as_os_str().is_empty() {
657 continue;
658 }
659 let dest = root.join(&rel);
660 if let Some(parent) = dest.parent() {
661 std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
662 }
663 let mut entry = entry;
664 let mut content = Vec::new();
665 entry.read_to_end(&mut content).map_err(|e| e.to_string())?;
666 std::fs::write(&dest, &content).map_err(|e| e.to_string())?;
667 total += content.len() as i64;
668 }
669 let _ = std::fs::create_dir_all(&root); // empty dir artifact still exists
670 Ok((total, true))
671}
672
673/// Build an uncompressed tar of the checkout, rooted at `workspace/` so it
674/// extracts to `/workspace` when uploaded to the container root.
675fn build_tar(files: &[TreeFile]) -> Vec<u8> {
676 let mut builder = tar::Builder::new(Vec::new());
677 for f in files {
678 let mut header = tar::Header::new_gnu();
679 header.set_size(f.content.len() as u64);
680 header.set_mode(if f.executable { 0o755 } else { 0o644 });
681 // append_data sets the path and checksum.
682 let _ = builder.append_data(
683 &mut header,
684 format!("workspace/{}", f.path),
685 f.content.as_slice(),
686 );
687 }
688 builder.into_inner().unwrap_or_default()
689}
690
691/// Single-quote a string for safe interpolation into a shell command.
692fn single_quote(s: &str) -> String {
693 format!("'{}'", s.replace('\'', "'\\''"))
694}
695
696#[cfg(test)]
697mod tests {
698 use super::*;
699
700 /// Build a tar the way docker's archive endpoint does: entries rooted at
701 /// the requested item's basename.
702 fn tar_of(entries: &[(&str, Option<&str>)]) -> Vec<u8> {
703 let mut b = tar::Builder::new(Vec::new());
704 for (path, content) in entries {
705 let mut h = tar::Header::new_gnu();
706 match content {
707 Some(c) => {
708 h.set_size(c.len() as u64);
709 h.set_mode(0o644);
710 h.set_entry_type(tar::EntryType::Regular);
711 b.append_data(&mut h, path, c.as_bytes()).unwrap();
712 }
713 None => {
714 h.set_size(0);
715 h.set_mode(0o755);
716 h.set_entry_type(tar::EntryType::Directory);
717 b.append_data(&mut h, path, std::io::empty()).unwrap();
718 }
719 }
720 }
721 b.into_inner().unwrap()
722 }
723
724 fn spec(name: &str, path: &str, browse: bool) -> ArtifactSpec {
725 ArtifactSpec {
726 name: name.into(),
727 path: path.into(),
728 browse,
729 has_meta: false,
730 }
731 }
732
733 #[test]
734 fn stores_a_file_artifact_as_is() {
735 let dir = tempfile::tempdir().unwrap();
736 let tar = tar_of(&[("anvild", Some("ELF..."))]);
737 let (size, is_dir) = store_artifact(
738 &spec("bin", "target/release/anvild", false),
739 &tar,
740 dir.path(),
741 )
742 .unwrap();
743 assert!(!is_dir);
744 assert_eq!(size, 6);
745 assert_eq!(std::fs::read(dir.path().join("bin")).unwrap(), b"ELF...");
746 }
747
748 #[test]
749 fn stores_a_directory_artifact_as_tar_gz() {
750 let dir = tempfile::tempdir().unwrap();
751 let tar = tar_of(&[("doc", None), ("doc/index.html", Some("<html>"))]);
752 let (size, is_dir) =
753 store_artifact(&spec("doc", "target/doc", false), &tar, dir.path()).unwrap();
754 assert!(is_dir);
755 let stored = dir.path().join("doc.tar.gz");
756 assert_eq!(size, stored.metadata().unwrap().len() as i64);
757 // Round-trips through gzip back to the original tar bytes.
758 let mut gz = flate2::read::GzDecoder::new(std::fs::File::open(&stored).unwrap());
759 let mut bytes = Vec::new();
760 gz.read_to_end(&mut bytes).unwrap();
761 assert_eq!(bytes, tar);
762 }
763
764 #[test]
765 fn extracts_a_browsable_directory_artifact() {
766 let dir = tempfile::tempdir().unwrap();
767 let tar = tar_of(&[
768 ("doc", None),
769 ("doc/index.html", Some("<html>")),
770 ("doc/sub", None),
771 ("doc/sub/page.html", Some("<p>")),
772 ]);
773 let (size, is_dir) =
774 store_artifact(&spec("doc", "target/doc", true), &tar, dir.path()).unwrap();
775 assert!(is_dir);
776 assert_eq!(size, 6 + 3);
777 let root = dir.path().join("doc");
778 assert_eq!(std::fs::read(root.join("index.html")).unwrap(), b"<html>");
779 assert_eq!(std::fs::read(root.join("sub/page.html")).unwrap(), b"<p>");
780 }
781}