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