| 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 | #[derive(Default)] |
| 69 | struct Inner { |
| 70 | leases: Mutex<HashMap<i64, Lease>>, |
| 71 | /// Runner name → (advertised platform, last seen). What a runner says it |
| 72 | /// is; not a credential and not checked — everyone holding the token is |
| 73 | /// one principal, so a runner that lies about its architecture is only |
| 74 | /// lying to itself about which jobs it gets. |
| 75 | runners: Mutex<HashMap<String, (String, Instant)>>, |
| 76 | wake: tokio::sync::Notify, |
| 77 | /// Serializes claim attempts. Reading the queue and marking a run |
| 78 | /// `running` are separate awaits, so two runners polling at once could |
| 79 | /// otherwise both be handed the same job. |
| 80 | claim: tokio::sync::Mutex<()>, |
| 81 | } |
| 82 | |
| 83 | impl Dispatch { |
| 84 | pub fn new() -> Self { |
| 85 | Self::default() |
| 86 | } |
| 87 | |
| 88 | /// Wake every parked claim request. Called when a run is enqueued. |
| 89 | /// |
| 90 | /// `notify_waiters` rather than `notify_one`: a runner that wakes and finds |
| 91 | /// the queue already emptied by another simply parks again, whereas |
| 92 | /// `notify_one` can hand the permit to a runner that is about to give up. |
| 93 | pub fn wake(&self) { |
| 94 | self.inner.wake.notify_waiters(); |
| 95 | } |
| 96 | |
| 97 | /// Park until there might be work, or until `timeout` elapses. |
| 98 | pub async fn wait_for_work(&self, timeout: Duration) { |
| 99 | let _ = tokio::time::timeout(timeout, self.inner.wake.notified()).await; |
| 100 | } |
| 101 | |
| 102 | /// Hold for the whole of a claim attempt, so the queue read and the |
| 103 | /// `mark_running` that follows it are atomic with respect to other |
| 104 | /// runners. |
| 105 | pub async fn claim_guard(&self) -> tokio::sync::MutexGuard<'_, ()> { |
| 106 | self.inner.claim.lock().await |
| 107 | } |
| 108 | |
| 109 | /// Record that `runner` holds `run_id`, stashing the job's secrets for |
| 110 | /// masking when the result arrives. |
| 111 | pub fn claim(&self, run_id: i64, runner: &str, secrets: Vec<(String, String)>) { |
| 112 | self.inner.leases.lock().unwrap().insert( |
| 113 | run_id, |
| 114 | Lease { |
| 115 | runner: runner.to_string(), |
| 116 | expires_at: Instant::now() + LEASE_TTL, |
| 117 | secrets, |
| 118 | }, |
| 119 | ); |
| 120 | } |
| 121 | |
| 122 | /// Note that `runner` is alive and what platform it says it is. Called on |
| 123 | /// every claim and heartbeat, which is what keeps [`platforms`] honest |
| 124 | /// about a busy runner as well as an idle one. |
| 125 | /// |
| 126 | /// [`platforms`]: Dispatch::platforms |
| 127 | pub fn seen(&self, runner: &str, platform: &str) { |
| 128 | if platform.is_empty() { |
| 129 | return; |
| 130 | } |
| 131 | self.inner |
| 132 | .runners |
| 133 | .lock() |
| 134 | .unwrap() |
| 135 | .insert(runner.to_string(), (platform.to_string(), Instant::now())); |
| 136 | } |
| 137 | |
| 138 | /// Platforms with a runner behind them right now (within [`RUNNER_TTL`]). |
| 139 | /// |
| 140 | /// The dispatcher's answer to "is there anyone who could run this |
| 141 | /// natively?" — if not, an emulating runner may take the job rather than |
| 142 | /// leaving it queued forever. |
| 143 | pub fn platforms(&self) -> Vec<String> { |
| 144 | let now = Instant::now(); |
| 145 | let mut runners = self.inner.runners.lock().unwrap(); |
| 146 | runners.retain(|_, (_, seen)| now.duration_since(*seen) < RUNNER_TTL); |
| 147 | let mut platforms: Vec<String> = runners.values().map(|(p, _)| p.clone()).collect(); |
| 148 | platforms.sort(); |
| 149 | platforms.dedup(); |
| 150 | platforms |
| 151 | } |
| 152 | |
| 153 | /// Every runner seen within [`RUNNER_TTL`], as (name, platform), for the |
| 154 | /// admin view and for logging. |
| 155 | pub fn runners(&self) -> Vec<(String, String)> { |
| 156 | let now = Instant::now(); |
| 157 | let mut runners = self.inner.runners.lock().unwrap(); |
| 158 | runners.retain(|_, (_, seen)| now.duration_since(*seen) < RUNNER_TTL); |
| 159 | let mut out: Vec<(String, String)> = runners |
| 160 | .iter() |
| 161 | .map(|(name, (platform, _))| (name.clone(), platform.clone())) |
| 162 | .collect(); |
| 163 | out.sort(); |
| 164 | out |
| 165 | } |
| 166 | |
| 167 | /// The secrets handed to `run_id`, for masking its log. |
| 168 | pub fn secrets_for(&self, run_id: i64) -> Vec<(String, String)> { |
| 169 | self.inner |
| 170 | .leases |
| 171 | .lock() |
| 172 | .unwrap() |
| 173 | .get(&run_id) |
| 174 | .map(|l| l.secrets.clone()) |
| 175 | .unwrap_or_default() |
| 176 | } |
| 177 | |
| 178 | /// Extend `runner`'s lease on `run_id`. False when it holds no such lease, |
| 179 | /// which is how a runner learns its job was requeued out from under it. |
| 180 | pub fn heartbeat(&self, run_id: i64, runner: &str) -> bool { |
| 181 | let mut leases = self.inner.leases.lock().unwrap(); |
| 182 | match leases.get_mut(&run_id) { |
| 183 | Some(lease) if lease.runner == runner => { |
| 184 | lease.expires_at = Instant::now() + LEASE_TTL; |
| 185 | true |
| 186 | } |
| 187 | _ => false, |
| 188 | } |
| 189 | } |
| 190 | |
| 191 | /// Whether `runner` currently holds `run_id`. The authorization check for |
| 192 | /// every per-job endpoint: a valid token gets you a job, but only the |
| 193 | /// holder may fetch its checkout or post its result. |
| 194 | pub fn holds(&self, run_id: i64, runner: &str) -> bool { |
| 195 | self.inner |
| 196 | .leases |
| 197 | .lock() |
| 198 | .unwrap() |
| 199 | .get(&run_id) |
| 200 | .is_some_and(|l| l.runner == runner) |
| 201 | } |
| 202 | |
| 203 | /// Whether anyone holds `run_id`. Lets the claim loop skip a run another |
| 204 | /// runner took between the queue read and here. |
| 205 | pub fn holds_any(&self, run_id: i64) -> bool { |
| 206 | self.inner.leases.lock().unwrap().contains_key(&run_id) |
| 207 | } |
| 208 | |
| 209 | /// Drop the lease on `run_id` (the job finished, one way or another). |
| 210 | pub fn release(&self, run_id: i64) { |
| 211 | self.inner.leases.lock().unwrap().remove(&run_id); |
| 212 | } |
| 213 | |
| 214 | /// Remove and return every lease past its TTL. The caller requeues them. |
| 215 | pub fn take_expired(&self) -> Vec<(i64, String)> { |
| 216 | let now = Instant::now(); |
| 217 | let mut leases = self.inner.leases.lock().unwrap(); |
| 218 | let dead: Vec<_> = leases |
| 219 | .iter() |
| 220 | .filter(|(_, l)| l.expires_at <= now) |
| 221 | .map(|(id, l)| (*id, l.runner.clone())) |
| 222 | .collect(); |
| 223 | for (id, _) in &dead { |
| 224 | leases.remove(id); |
| 225 | } |
| 226 | dead |
| 227 | } |
| 228 | |
| 229 | /// Run ids currently claimed, for the admin view and for logging. |
| 230 | pub fn in_flight(&self) -> Vec<(i64, String)> { |
| 231 | self.inner |
| 232 | .leases |
| 233 | .lock() |
| 234 | .unwrap() |
| 235 | .iter() |
| 236 | .map(|(id, l)| (*id, l.runner.clone())) |
| 237 | .collect() |
| 238 | } |
| 239 | } |
| 240 | |
| 241 | #[cfg(test)] |
| 242 | mod tests { |
| 243 | use super::*; |
| 244 | |
| 245 | /// The dispatcher's view of who is out there: one entry per runner, one |
| 246 | /// platform per architecture, and a re-registration under the same name |
| 247 | /// (a runner restarted on a rebuilt machine) replaces rather than doubles. |
| 248 | #[test] |
| 249 | fn platforms_are_deduped_per_live_runner() { |
| 250 | let jobs = Dispatch::new(); |
| 251 | assert!(jobs.platforms().is_empty()); |
| 252 | |
| 253 | jobs.seen("mac", "linux/arm64"); |
| 254 | jobs.seen("droplet", "linux/amd64"); |
| 255 | jobs.seen("laptop", "linux/arm64"); |
| 256 | assert_eq!(jobs.platforms(), vec!["linux/amd64", "linux/arm64"]); |
| 257 | assert_eq!(jobs.runners().len(), 3); |
| 258 | |
| 259 | jobs.seen("mac", "linux/amd64"); // same name, rebuilt as amd64 |
| 260 | assert_eq!(jobs.runners().len(), 3); |
| 261 | assert_eq!(jobs.platforms(), vec!["linux/amd64", "linux/arm64"]); |
| 262 | |
| 263 | // A runner that advertises nothing is not evidence of any platform, |
| 264 | // and must not register as one — routing would then never fall back. |
| 265 | jobs.seen("mystery", ""); |
| 266 | assert_eq!(jobs.runners().len(), 3); |
| 267 | } |
| 268 | } |