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/// Shared dispatch state. Cheap to clone; all clones share one map.
53#[derive(Clone, Default)]
54pub struct Dispatch {
55 inner: Arc<Inner>,
56}
57
58#[derive(Default)]
59struct 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
68impl 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}