| 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 | /// Shared dispatch state. Cheap to clone; all clones share one map. |
| 53 | #[derive(Clone, Default)] |
| 54 | pub struct Dispatch { |
| 55 | inner: Arc<Inner>, |
| 56 | } |
| 57 | |
| 58 | #[derive(Default)] |
| 59 | struct Inner { |
| 60 | leases: Mutex<HashMap<i64, Lease>>, |
| 61 | wake: tokio::sync::Notify, |
| 62 | /// Serializes claim attempts. Reading the queue and marking a run |
| 63 | /// `running` are separate awaits, so two runners polling at once could |
| 64 | /// otherwise both be handed the same job. |
| 65 | claim: tokio::sync::Mutex<()>, |
| 66 | } |
| 67 | |
| 68 | impl Dispatch { |
| 69 | pub fn new() -> Self { |
| 70 | Self::default() |
| 71 | } |
| 72 | |
| 73 | /// Wake every parked claim request. Called when a run is enqueued. |
| 74 | /// |
| 75 | /// `notify_waiters` rather than `notify_one`: a runner that wakes and finds |
| 76 | /// the queue already emptied by another simply parks again, whereas |
| 77 | /// `notify_one` can hand the permit to a runner that is about to give up. |
| 78 | pub fn wake(&self) { |
| 79 | self.inner.wake.notify_waiters(); |
| 80 | } |
| 81 | |
| 82 | /// Park until there might be work, or until `timeout` elapses. |
| 83 | pub async fn wait_for_work(&self, timeout: Duration) { |
| 84 | let _ = tokio::time::timeout(timeout, self.inner.wake.notified()).await; |
| 85 | } |
| 86 | |
| 87 | /// Hold for the whole of a claim attempt, so the queue read and the |
| 88 | /// `mark_running` that follows it are atomic with respect to other |
| 89 | /// runners. |
| 90 | pub async fn claim_guard(&self) -> tokio::sync::MutexGuard<'_, ()> { |
| 91 | self.inner.claim.lock().await |
| 92 | } |
| 93 | |
| 94 | /// Record that `runner` holds `run_id`, stashing the job's secrets for |
| 95 | /// masking when the result arrives. |
| 96 | pub fn claim(&self, run_id: i64, runner: &str, secrets: Vec<(String, String)>) { |
| 97 | self.inner.leases.lock().unwrap().insert( |
| 98 | run_id, |
| 99 | Lease { |
| 100 | runner: runner.to_string(), |
| 101 | expires_at: Instant::now() + LEASE_TTL, |
| 102 | secrets, |
| 103 | }, |
| 104 | ); |
| 105 | } |
| 106 | |
| 107 | /// The secrets handed to `run_id`, for masking its log. |
| 108 | pub fn secrets_for(&self, run_id: i64) -> Vec<(String, String)> { |
| 109 | self.inner |
| 110 | .leases |
| 111 | .lock() |
| 112 | .unwrap() |
| 113 | .get(&run_id) |
| 114 | .map(|l| l.secrets.clone()) |
| 115 | .unwrap_or_default() |
| 116 | } |
| 117 | |
| 118 | /// Extend `runner`'s lease on `run_id`. False when it holds no such lease, |
| 119 | /// which is how a runner learns its job was requeued out from under it. |
| 120 | pub fn heartbeat(&self, run_id: i64, runner: &str) -> bool { |
| 121 | let mut leases = self.inner.leases.lock().unwrap(); |
| 122 | match leases.get_mut(&run_id) { |
| 123 | Some(lease) if lease.runner == runner => { |
| 124 | lease.expires_at = Instant::now() + LEASE_TTL; |
| 125 | true |
| 126 | } |
| 127 | _ => false, |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | /// Whether `runner` currently holds `run_id`. The authorization check for |
| 132 | /// every per-job endpoint: a valid token gets you a job, but only the |
| 133 | /// holder may fetch its checkout or post its result. |
| 134 | pub fn holds(&self, run_id: i64, runner: &str) -> bool { |
| 135 | self.inner |
| 136 | .leases |
| 137 | .lock() |
| 138 | .unwrap() |
| 139 | .get(&run_id) |
| 140 | .is_some_and(|l| l.runner == runner) |
| 141 | } |
| 142 | |
| 143 | /// Whether anyone holds `run_id`. Lets the claim loop skip a run another |
| 144 | /// runner took between the queue read and here. |
| 145 | pub fn holds_any(&self, run_id: i64) -> bool { |
| 146 | self.inner.leases.lock().unwrap().contains_key(&run_id) |
| 147 | } |
| 148 | |
| 149 | /// Drop the lease on `run_id` (the job finished, one way or another). |
| 150 | pub fn release(&self, run_id: i64) { |
| 151 | self.inner.leases.lock().unwrap().remove(&run_id); |
| 152 | } |
| 153 | |
| 154 | /// Remove and return every lease past its TTL. The caller requeues them. |
| 155 | pub fn take_expired(&self) -> Vec<(i64, String)> { |
| 156 | let now = Instant::now(); |
| 157 | let mut leases = self.inner.leases.lock().unwrap(); |
| 158 | let dead: Vec<_> = leases |
| 159 | .iter() |
| 160 | .filter(|(_, l)| l.expires_at <= now) |
| 161 | .map(|(id, l)| (*id, l.runner.clone())) |
| 162 | .collect(); |
| 163 | for (id, _) in &dead { |
| 164 | leases.remove(id); |
| 165 | } |
| 166 | dead |
| 167 | } |
| 168 | |
| 169 | /// Run ids currently claimed, for the admin view and for logging. |
| 170 | pub fn in_flight(&self) -> Vec<(i64, String)> { |
| 171 | self.inner |
| 172 | .leases |
| 173 | .lock() |
| 174 | .unwrap() |
| 175 | .iter() |
| 176 | .map(|(id, l)| (*id, l.runner.clone())) |
| 177 | .collect() |
| 178 | } |
| 179 | } |