| 1 | //! Dispatch state for remote job runners (see `docs/remote-runners.md`). |
| 2 | //! |
| 3 | //! anvil no longer executes CI itself; it hands jobs to runners that dial in |
| 4 | //! and claim them. This holds the two pieces of state that implies: who is |
| 5 | //! currently running what, and a way to wake a parked claim request when work |
| 6 | //! arrives. |
| 7 | //! |
| 8 | //! **Deliberately in memory**, not columns on `CiRun`. Toasty has no migration |
| 9 | //! support yet (see `TODO.md` and DEPLOY.md's Operations note), so adding |
| 10 | //! `claimed_by`/`lease_expires_at` to the model would need a hand-written |
| 11 | //! migration against the live database. It also costs nothing to keep here: |
| 12 | //! anvild is the only dispatcher, so a lease has no reason to outlive it, and |
| 13 | //! the case where it dies mid-job is already covered — `ci::requeue_interrupted` |
| 14 | //! re-queues everything still marked `running` at startup. |
| 15 | //! |
| 16 | //! Same shape as [`crate::agent::Registry`], for the same reason. |
| 17 | |
| 18 | use std::{ |
| 19 | collections::HashMap, |
| 20 | sync::{ |
| 21 | Arc, |
| 22 | Mutex, |
| 23 | }, |
| 24 | time::{ |
| 25 | Duration, |
| 26 | Instant, |
| 27 | }, |
| 28 | }; |
| 29 | |
| 30 | pub use anvil_job::{ |
| 31 | HEARTBEAT_INTERVAL, |
| 32 | LEASE_TTL, |
| 33 | }; |
| 34 | |
| 35 | /// Who holds a run, and until when. |
| 36 | #[derive(Clone, Debug)] |
| 37 | pub struct Lease { |
| 38 | pub runner: String, |
| 39 | pub expires_at: Instant, |
| 40 | /// The secret values handed to this job, kept so the log can be masked |
| 41 | /// when the result comes back. |
| 42 | /// |
| 43 | /// Stashed rather than re-read from the vault at finish time, because |
| 44 | /// `Vault::take` fails once the repository's unlock TTL lapses — and a job |
| 45 | /// may well outlive it. Re-reading would mean a run that took longer than |
| 46 | /// the unlock silently skips masking, which is exactly the case where the |
| 47 | /// log is most likely to contain something. No new exposure: these values |
| 48 | /// are already in this process's vault. |
| 49 | pub secrets: Vec<(String, String)>, |
| 50 | } |
| 51 | |
| 52 | /// How long a runner counts as present after its last claim or heartbeat. |
| 53 | /// |
| 54 | /// Only platform routing reads this: a job that wants `linux/amd64` waits for |
| 55 | /// a native runner while one is present, and is handed to an emulating runner |
| 56 | /// once none is. Comfortably longer than both the claim poll (55s, so an idle |
| 57 | /// runner refreshes itself) and [`HEARTBEAT_INTERVAL`] (30s, so a runner stays |
| 58 | /// present through a long build), and short enough that a machine that went to |
| 59 | /// sleep stops holding its architecture's jobs hostage for long. |
| 60 | pub const RUNNER_TTL: Duration = Duration::from_secs(300); |
| 61 | |
| 62 | /// Shared dispatch state. Cheap to clone; all clones share one map. |
| 63 | #[derive(Clone, Default)] |
| 64 | pub struct Dispatch { |
| 65 | inner: Arc<Inner>, |
| 66 | } |
| 67 | |
| 68 | /// What a runner has told us about itself, and when it last did. |
| 69 | /// |
| 70 | /// The platform and version are self-asserted and not checked — everyone |
| 71 | /// holding the token is one principal, so a runner that lies about its |
| 72 | /// architecture is only lying to itself about which jobs it gets. |
| 73 | struct Presence { |
| 74 | platform: String, |
| 75 | version: String, |
| 76 | /// First contact since this process started. Not the runner's own uptime: |
| 77 | /// an anvild restart resets it, and the page says as much. |
| 78 | first_seen: Instant, |
| 79 | last_seen: Instant, |
| 80 | } |
| 81 | |
| 82 | /// One runner as [`Dispatch::runners`] reports it: what it says it is, how |
| 83 | /// stale that is, and what it is holding right now. |
| 84 | #[derive(Clone, Debug)] |
| 85 | pub struct RunnerStatus { |
| 86 | pub name: String, |
| 87 | /// Empty if the runner has only ever reached endpoints that identify it by |
| 88 | /// header (no JSON body to carry a platform). |
| 89 | pub platform: String, |
| 90 | pub version: String, |
| 91 | /// Since its last claim or heartbeat. An idle runner refreshes this once |
| 92 | /// per claim poll (55s) and a busy one every [`HEARTBEAT_INTERVAL`], so |
| 93 | /// anything much past a minute means it stopped talking. |
| 94 | pub last_seen: Duration, |
| 95 | /// Since anvild first heard from it, reset by an anvild restart. |
| 96 | pub connected_for: Duration, |
| 97 | /// Runs it currently holds a lease on. |
| 98 | pub running: Vec<i64>, |
| 99 | } |
| 100 | |
| 101 | #[derive(Default)] |
| 102 | struct Inner { |
| 103 | leases: Mutex<HashMap<i64, Lease>>, |
| 104 | /// Runner name → what it last told us. Presence is derived from ordinary |
| 105 | /// traffic (claims and heartbeats) rather than a dedicated ping: a runner |
| 106 | /// that is not polling for work or beating for a job is not a runner |
| 107 | /// anything should be dispatched to, whatever a ping would claim. |
| 108 | runners: Mutex<HashMap<String, Presence>>, |
| 109 | wake: tokio::sync::Notify, |
| 110 | /// Serializes claim attempts. Reading the queue and marking a run |
| 111 | /// `running` are separate awaits, so two runners polling at once could |
| 112 | /// otherwise both be handed the same job. |
| 113 | claim: tokio::sync::Mutex<()>, |
| 114 | } |
| 115 | |
| 116 | impl Dispatch { |
| 117 | pub fn new() -> Self { |
| 118 | Self::default() |
| 119 | } |
| 120 | |
| 121 | /// Wake every parked claim request. Called when a run is enqueued. |
| 122 | /// |
| 123 | /// `notify_waiters` rather than `notify_one`: a runner that wakes and finds |
| 124 | /// the queue already emptied by another simply parks again, whereas |
| 125 | /// `notify_one` can hand the permit to a runner that is about to give up. |
| 126 | pub fn wake(&self) { |
| 127 | self.inner.wake.notify_waiters(); |
| 128 | } |
| 129 | |
| 130 | /// Park until there might be work, or until `timeout` elapses. |
| 131 | pub async fn wait_for_work(&self, timeout: Duration) { |
| 132 | let _ = tokio::time::timeout(timeout, self.inner.wake.notified()).await; |
| 133 | } |
| 134 | |
| 135 | /// Hold for the whole of a claim attempt, so the queue read and the |
| 136 | /// `mark_running` that follows it are atomic with respect to other |
| 137 | /// runners. |
| 138 | pub async fn claim_guard(&self) -> tokio::sync::MutexGuard<'_, ()> { |
| 139 | self.inner.claim.lock().await |
| 140 | } |
| 141 | |
| 142 | /// Record that `runner` holds `run_id`, stashing the job's secrets for |
| 143 | /// masking when the result arrives. |
| 144 | pub fn claim(&self, run_id: i64, runner: &str, secrets: Vec<(String, String)>) { |
| 145 | self.inner.leases.lock().unwrap().insert( |
| 146 | run_id, |
| 147 | Lease { |
| 148 | runner: runner.to_string(), |
| 149 | expires_at: Instant::now() + LEASE_TTL, |
| 150 | secrets, |
| 151 | }, |
| 152 | ); |
| 153 | } |
| 154 | |
| 155 | /// Note that `runner` is alive, and what it says it is. Called on every |
| 156 | /// claim and heartbeat, which is what keeps [`platforms`] honest about a |
| 157 | /// busy runner as well as an idle one. |
| 158 | /// |
| 159 | /// An empty `platform` or `version` means "not stated on this request" |
| 160 | /// rather than "unknown": the endpoints with no JSON body identify their |
| 161 | /// caller by header alone, and must not blank out what a claim established. |
| 162 | /// |
| 163 | /// [`platforms`]: Dispatch::platforms |
| 164 | pub fn seen(&self, runner: &str, platform: &str, version: &str) { |
| 165 | let now = Instant::now(); |
| 166 | let mut runners = self.inner.runners.lock().unwrap(); |
| 167 | let entry = runners |
| 168 | .entry(runner.to_string()) |
| 169 | .or_insert_with(|| Presence { |
| 170 | platform: String::new(), |
| 171 | version: String::new(), |
| 172 | first_seen: now, |
| 173 | last_seen: now, |
| 174 | }); |
| 175 | if !platform.is_empty() { |
| 176 | entry.platform = platform.to_string(); |
| 177 | } |
| 178 | if !version.is_empty() { |
| 179 | entry.version = version.to_string(); |
| 180 | } |
| 181 | entry.last_seen = now; |
| 182 | } |
| 183 | |
| 184 | /// Drop everyone past [`RUNNER_TTL`] and hand back the rest. Pruned on read |
| 185 | /// rather than on a timer: presence is only ever consulted here. |
| 186 | fn live(&self) -> std::sync::MutexGuard<'_, HashMap<String, Presence>> { |
| 187 | let now = Instant::now(); |
| 188 | let mut runners = self.inner.runners.lock().unwrap(); |
| 189 | runners.retain(|_, p| now.duration_since(p.last_seen) < RUNNER_TTL); |
| 190 | runners |
| 191 | } |
| 192 | |
| 193 | /// Platforms with a runner behind them right now (within [`RUNNER_TTL`]). |
| 194 | /// |
| 195 | /// The dispatcher's answer to "is there anyone who could run this |
| 196 | /// natively?" — if not, an emulating runner may take the job rather than |
| 197 | /// leaving it queued forever. A runner that has stated no platform is |
| 198 | /// present but is not evidence of any architecture, so it is skipped here: |
| 199 | /// counting it would strand every job behind a fallback that never fires. |
| 200 | pub fn platforms(&self) -> Vec<String> { |
| 201 | let mut platforms: Vec<String> = self |
| 202 | .live() |
| 203 | .values() |
| 204 | .filter(|p| !p.platform.is_empty()) |
| 205 | .map(|p| p.platform.clone()) |
| 206 | .collect(); |
| 207 | platforms.sort(); |
| 208 | platforms.dedup(); |
| 209 | platforms |
| 210 | } |
| 211 | |
| 212 | /// Every runner seen within [`RUNNER_TTL`], for the admin view and for |
| 213 | /// logging. Sorted by name, so the page does not reshuffle between polls. |
| 214 | pub fn runners(&self) -> Vec<RunnerStatus> { |
| 215 | // Leases first and released before the presence lock is taken: the two |
| 216 | // are never held together anywhere, which is what keeps them from |
| 217 | // needing an ordering rule. |
| 218 | let mut running: HashMap<String, Vec<i64>> = HashMap::new(); |
| 219 | for (run_id, runner) in self.in_flight() { |
| 220 | running.entry(runner).or_default().push(run_id); |
| 221 | } |
| 222 | |
| 223 | let now = Instant::now(); |
| 224 | let mut out: Vec<RunnerStatus> = self |
| 225 | .live() |
| 226 | .iter() |
| 227 | .map(|(name, p)| RunnerStatus { |
| 228 | name: name.clone(), |
| 229 | platform: p.platform.clone(), |
| 230 | version: p.version.clone(), |
| 231 | last_seen: now.duration_since(p.last_seen), |
| 232 | connected_for: now.duration_since(p.first_seen), |
| 233 | running: { |
| 234 | let mut ids = running.remove(name).unwrap_or_default(); |
| 235 | ids.sort_unstable(); |
| 236 | ids |
| 237 | }, |
| 238 | }) |
| 239 | .collect(); |
| 240 | out.sort_by(|a, b| a.name.cmp(&b.name)); |
| 241 | out |
| 242 | } |
| 243 | |
| 244 | /// The secrets handed to `run_id`, for masking its log. |
| 245 | pub fn secrets_for(&self, run_id: i64) -> Vec<(String, String)> { |
| 246 | self.inner |
| 247 | .leases |
| 248 | .lock() |
| 249 | .unwrap() |
| 250 | .get(&run_id) |
| 251 | .map(|l| l.secrets.clone()) |
| 252 | .unwrap_or_default() |
| 253 | } |
| 254 | |
| 255 | /// Extend `runner`'s lease on `run_id`. False when it holds no such lease, |
| 256 | /// which is how a runner learns its job was requeued out from under it. |
| 257 | pub fn heartbeat(&self, run_id: i64, runner: &str) -> bool { |
| 258 | let mut leases = self.inner.leases.lock().unwrap(); |
| 259 | match leases.get_mut(&run_id) { |
| 260 | Some(lease) if lease.runner == runner => { |
| 261 | lease.expires_at = Instant::now() + LEASE_TTL; |
| 262 | true |
| 263 | } |
| 264 | _ => false, |
| 265 | } |
| 266 | } |
| 267 | |
| 268 | /// Whether `runner` currently holds `run_id`. The authorization check for |
| 269 | /// every per-job endpoint: a valid token gets you a job, but only the |
| 270 | /// holder may fetch its checkout or post its result. |
| 271 | pub fn holds(&self, run_id: i64, runner: &str) -> bool { |
| 272 | self.inner |
| 273 | .leases |
| 274 | .lock() |
| 275 | .unwrap() |
| 276 | .get(&run_id) |
| 277 | .is_some_and(|l| l.runner == runner) |
| 278 | } |
| 279 | |
| 280 | /// Whether anyone holds `run_id`. Lets the claim loop skip a run another |
| 281 | /// runner took between the queue read and here. |
| 282 | pub fn holds_any(&self, run_id: i64) -> bool { |
| 283 | self.inner.leases.lock().unwrap().contains_key(&run_id) |
| 284 | } |
| 285 | |
| 286 | /// Drop the lease on `run_id` (the job finished, one way or another). |
| 287 | pub fn release(&self, run_id: i64) { |
| 288 | self.inner.leases.lock().unwrap().remove(&run_id); |
| 289 | } |
| 290 | |
| 291 | /// Remove and return every lease past its TTL. The caller requeues them. |
| 292 | pub fn take_expired(&self) -> Vec<(i64, String)> { |
| 293 | let now = Instant::now(); |
| 294 | let mut leases = self.inner.leases.lock().unwrap(); |
| 295 | let dead: Vec<_> = leases |
| 296 | .iter() |
| 297 | .filter(|(_, l)| l.expires_at <= now) |
| 298 | .map(|(id, l)| (*id, l.runner.clone())) |
| 299 | .collect(); |
| 300 | for (id, _) in &dead { |
| 301 | leases.remove(id); |
| 302 | } |
| 303 | dead |
| 304 | } |
| 305 | |
| 306 | /// Run ids currently claimed, for the admin view and for logging. |
| 307 | pub fn in_flight(&self) -> Vec<(i64, String)> { |
| 308 | self.inner |
| 309 | .leases |
| 310 | .lock() |
| 311 | .unwrap() |
| 312 | .iter() |
| 313 | .map(|(id, l)| (*id, l.runner.clone())) |
| 314 | .collect() |
| 315 | } |
| 316 | } |
| 317 | |
| 318 | #[cfg(test)] |
| 319 | mod tests { |
| 320 | use super::*; |
| 321 | |
| 322 | /// The dispatcher's view of who is out there: one entry per runner, one |
| 323 | /// platform per architecture, and a re-registration under the same name |
| 324 | /// (a runner restarted on a rebuilt machine) replaces rather than doubles. |
| 325 | #[test] |
| 326 | fn platforms_are_deduped_per_live_runner() { |
| 327 | let jobs = Dispatch::new(); |
| 328 | assert!(jobs.platforms().is_empty()); |
| 329 | |
| 330 | jobs.seen("mac", "linux/arm64", "0.1.0"); |
| 331 | jobs.seen("droplet", "linux/amd64", "0.1.0"); |
| 332 | jobs.seen("laptop", "linux/arm64", "0.1.0"); |
| 333 | assert_eq!(jobs.platforms(), vec!["linux/amd64", "linux/arm64"]); |
| 334 | assert_eq!(jobs.runners().len(), 3); |
| 335 | |
| 336 | jobs.seen("mac", "linux/amd64", "0.1.0"); // same name, rebuilt as amd64 |
| 337 | assert_eq!(jobs.runners().len(), 3); |
| 338 | assert_eq!(jobs.platforms(), vec!["linux/amd64", "linux/arm64"]); |
| 339 | |
| 340 | // A runner that advertises nothing is present — the status page should |
| 341 | // show it — but is not evidence of any platform, and must not register |
| 342 | // as one: routing would then never fall back to an emulating runner. |
| 343 | jobs.seen("mystery", "", ""); |
| 344 | assert_eq!(jobs.runners().len(), 4); |
| 345 | assert_eq!(jobs.platforms(), vec!["linux/amd64", "linux/arm64"]); |
| 346 | |
| 347 | // A later request that states nothing (the checkout GET and the |
| 348 | // artifact/result POSTs identify their caller by header alone) is a |
| 349 | // liveness signal, not an erasure of what the claim established. |
| 350 | jobs.seen("mac", "", ""); |
| 351 | let mac = jobs |
| 352 | .runners() |
| 353 | .into_iter() |
| 354 | .find(|r| r.name == "mac") |
| 355 | .expect("mac is present"); |
| 356 | assert_eq!(mac.platform, "linux/amd64"); |
| 357 | assert_eq!(mac.version, "0.1.0"); |
| 358 | } |
| 359 | |
| 360 | /// The status page's other column: who is holding a run right now. |
| 361 | #[test] |
| 362 | fn runners_report_the_leases_they_hold() { |
| 363 | let jobs = Dispatch::new(); |
| 364 | jobs.seen("mac", "linux/arm64", "0.1.0"); |
| 365 | jobs.seen("droplet", "linux/amd64", "0.1.0"); |
| 366 | jobs.claim(7, "mac", vec![]); |
| 367 | jobs.claim(9, "mac", vec![]); |
| 368 | |
| 369 | let by_name: HashMap<String, RunnerStatus> = jobs |
| 370 | .runners() |
| 371 | .into_iter() |
| 372 | .map(|r| (r.name.clone(), r)) |
| 373 | .collect(); |
| 374 | assert_eq!(by_name["mac"].running, vec![7, 9]); |
| 375 | assert!(by_name["droplet"].running.is_empty()); |
| 376 | |
| 377 | jobs.release(7); |
| 378 | jobs.release(9); |
| 379 | assert!(jobs.runners().iter().all(|r| r.running.is_empty())); |
| 380 | } |
| 381 | } |