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/// 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.
73struct 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)]
85pub 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)]
102struct 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
116impl 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)]
319mod 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}