| 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 | |
| 19 | mod client; |
| 20 | mod executor; |
| 21 | |
| 22 | use std::{ |
| 23 | sync::Arc, |
| 24 | time::Duration, |
| 25 | }; |
| 26 | |
| 27 | use anvil_job::{ |
| 28 | JobResult, |
| 29 | JobSpec, |
| 30 | RunnerInfo, |
| 31 | }; |
| 32 | use clap::Parser; |
| 33 | use 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. |
| 41 | const RETRY_DELAY: Duration = Duration::from_secs(15); |
| 42 | |
| 43 | #[derive(Parser)] |
| 44 | #[command(name = "anvil-worker", about = "Claim and run anvil CI jobs")] |
| 45 | struct 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] |
| 60 | async 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. |
| 115 | async 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 | |
| 170 | async 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. |
| 177 | async 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. |
| 189 | fn 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 | |
| 200 | fn hostname() -> String { |
| 201 | std::env::var("HOSTNAME") |
| 202 | .ok() |
| 203 | .filter(|h| !h.is_empty()) |
| 204 | .unwrap_or_else(|| "runner".to_string()) |
| 205 | } |