anvilsign in

collin/anvil

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
18use std::{
19 collections::HashMap,
20 sync::{
21 Arc,
22 Mutex,
23 },
24 time::{
25 Duration,
26 Instant,
27 },
28};
29
30pub use anvil_job::{
31 HEARTBEAT_INTERVAL,
32 LEASE_TTL,
33};
34
35/// Who holds a run, and until when.
36#[derive(Clone, Debug)]
37pub 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.
60pub 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)]
64pub struct Dispatch {
65 inner: Arc<Inner>,
66}
67
68#[derive(Default)]
69struct 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
83impl 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)]
242mod 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}