anvilsign in

collin/anvil

main / crates / anvil-worker / src / main.rs
1//! anvil's job runner.
2//!
3//! Claims CI jobs from a forge and runs them in sandboxed containers on this
4//! machine's Docker daemon. Runs one job at a time, forever, reconnecting
5//! through anything that goes wrong.
6//!
7//! ```sh
8//! anvil-worker --url https://anvil.example.com --token "$ANVIL_RUNNER_TOKEN"
9//! ```
10//!
11//! On macOS run it natively (a launchd agent), not in a container: a
12//! containerized runner would need the Docker socket mounted into it, which
13//! rebuilds exactly the root-equivalent hole moving execution off the forge
14//! was meant to remove. Isolation is unaffected by the host being a Mac —
15//! Docker Desktop runs every container inside one Linux VM, so the sandbox is
16//! enforced by the same kernel primitives as on Linux, with the VM as an extra
17//! boundary a bare-metal Linux host does not have.
18
19mod client;
20mod executor;
21
22use std::{
23 sync::Arc,
24 time::Duration,
25};
26
27use anvil_job::{
28 JobResult,
29 JobSpec,
30 RunnerInfo,
31};
32use clap::Parser;
33use client::{
34 Client,
35 UploadSink,
36};
37
38/// How long to wait before retrying after a failed claim. Long enough not to
39/// hammer a forge that is down or misconfigured; short enough that a runner
40/// picks up again promptly once it is back.
41const RETRY_DELAY: Duration = Duration::from_secs(15);
42
43#[derive(Parser)]
44#[command(name = "anvil-worker", about = "Claim and run anvil CI jobs")]
45struct Args {
46 /// Base URL of the anvil instance, e.g. https://anvil.example.com
47 #[arg(long, env = "ANVIL_URL")]
48 url: String,
49
50 /// Shared secret matching the forge's `[ci] runner_token`.
51 #[arg(long, env = "ANVIL_RUNNER_TOKEN", hide_env_values = true)]
52 token: String,
53
54 /// Name for this runner in logs and run headers. Defaults to the hostname.
55 #[arg(long, env = "ANVIL_RUNNER_NAME")]
56 name: Option<String>,
57}
58
59#[tokio::main]
60async fn main() {
61 tracing_subscriber::fmt()
62 .with_env_filter(
63 tracing_subscriber::EnvFilter::try_from_default_env()
64 .unwrap_or_else(|_| "anvil_worker=info".into()),
65 )
66 .init();
67
68 let args = Args::parse();
69 let name = args.name.unwrap_or_else(hostname);
70 let info = RunnerInfo {
71 name: name.clone(),
72 platform: native_platform(),
73 version: env!("CARGO_PKG_VERSION").to_string(),
74 };
75 tracing::info!(
76 "anvil-worker {} ({}) → {}",
77 info.name,
78 info.platform,
79 args.url
80 );
81
82 // Fail fast on a daemon that is not there, rather than claiming a job and
83 // immediately erroring it.
84 if let Err(e) = anvil_docker::connect() {
85 tracing::error!("{e}");
86 std::process::exit(1);
87 }
88
89 let client = match Client::new(&args.url, &args.token, info) {
90 Ok(c) => Arc::new(c),
91 Err(e) => {
92 tracing::error!("{e}");
93 std::process::exit(1);
94 }
95 };
96
97 loop {
98 match client.claim().await {
99 Ok(Some(job)) => run(&client, job).await,
100 // The poll expired with nothing queued. Straight back in.
101 Ok(None) => {}
102 Err(e) => {
103 tracing::error!("{e}");
104 tokio::time::sleep(RETRY_DELAY).await;
105 }
106 }
107 }
108}
109
110/// Run one claimed job and report it.
111///
112/// Every failure path still reports: a job anvil handed out and never hears
113/// about again sits `running` until its lease expires, which is a slow and
114/// confusing way to learn that an image name was wrong.
115async fn run(client: &Arc<Client>, job: JobSpec) {
116 let run_id = job.run_id;
117 tracing::info!("run {run_id}: claimed ({})", job.image);
118
119 let tar = match client.checkout(run_id).await {
120 Ok(tar) => tar,
121 Err(e) => {
122 report(
123 client,
124 run_id,
125 JobResult {
126 exit_code: 0,
127 log: String::new(),
128 artifacts: Vec::new(),
129 runner_error: Some(e),
130 },
131 )
132 .await;
133 return;
134 }
135 };
136
137 // Heartbeat for as long as the job runs; aborted below once it is done.
138 let beat = tokio::spawn(heartbeat(Arc::clone(client), run_id));
139
140 let mut log = String::new();
141 let mut sink = UploadSink { client, run_id };
142 let result = match executor::execute(&job, tar, &mut log, &mut sink).await {
143 Ok((exit_code, artifacts)) => JobResult {
144 exit_code,
145 log,
146 artifacts,
147 runner_error: None,
148 },
149 Err(e) => JobResult {
150 exit_code: 0,
151 log,
152 artifacts: Vec::new(),
153 runner_error: Some(e),
154 },
155 };
156 beat.abort();
157
158 tracing::info!(
159 "run {run_id}: finished (exit {}{})",
160 result.exit_code,
161 result
162 .runner_error
163 .as_ref()
164 .map(|e| format!(", runner error: {e}"))
165 .unwrap_or_default()
166 );
167 report(client, run_id, result).await;
168}
169
170async fn report(client: &Client, run_id: i64, result: JobResult) {
171 if let Err(e) = client.report(run_id, &result).await {
172 tracing::error!("run {run_id}: reporting failed: {e}");
173 }
174}
175
176/// Keep the lease alive while a job runs.
177async fn heartbeat(client: Arc<Client>, run_id: i64) {
178 let mut tick = tokio::time::interval(anvil_job::HEARTBEAT_INTERVAL);
179 loop {
180 tick.tick().await;
181 if !client.heartbeat(run_id).await {
182 tracing::warn!("run {run_id}: anvil no longer thinks we hold this job");
183 }
184 }
185}
186
187/// `os/arch` in Docker's spelling, advertised so anvil can eventually route a
188/// job that needs a particular architecture to a runner that has it natively.
189fn native_platform() -> String {
190 let arch = match std::env::consts::ARCH {
191 "x86_64" => "amd64",
192 "aarch64" => "arm64",
193 other => other,
194 };
195 // Always `linux`: containers run in a Linux VM on macOS, so the daemon's
196 // platform is linux there too, whatever the host OS is.
197 format!("linux/{arch}")
198}
199
200fn hostname() -> String {
201 std::env::var("HOSTNAME")
202 .ok()
203 .filter(|h| !h.is_empty())
204 .unwrap_or_else(|| "runner".to_string())
205}