Skip to content
1,620 linesCodeBlameRaw
1//! Self-hosted runners: registering them, handing them work, hearing back,
2//! and the people's side (listing, removing, groups, settings). See
3//! `g1t_contracts::runners` for the model and the wire.
4//!
5//! Work reaches a runner only when it asks for it (`runner_poll`), so a
6//! runner never needs an open port. A workflow job is taken by flipping its
7//! row from `queued` to `in_progress` in one statement, exactly as g1t's
8//! own sandboxes take one, and the harness on the machine then fetches it
9//! and reports with the job's own token through the same API. Agent work
10//! from the runner service is kept as a task with its environment sealed,
11//! and handed over once.
12
13use g1t_contracts::access::Capability;
14use g1t_contracts::audit::{AuditActor, AuditOutcome, AuditTarget, NewAuditEntry, RecordAuditArgs, Surface};
15use g1t_contracts::billing::{ComputeKind, RecordSandboxArgs};
16use g1t_contracts::repos::{Repo, RepoPath};
17use g1t_contracts::runners::{
18 self as model, StuckJob, StuckJobsArgs, Assignment, CancelTaskArgs, CreateRegistrationTokenArgs, DeleteRunnerGroupArgs, EnqueueTaskArgs, FinishedArgs,
19 ListRunnersArgs, Poll, PollArgs, RegisterArgs, Registered, RegistrationToken, RemoveRunnerArgs, RemoveSelfArgs, RouteArgs, Runner,
20 RunnerAuth, RunnerGroup, RunnerGroupsArgs, RunnerSettings, RunnerSettingsArgs, RunnerWork, RunnersOwner, SetRunnerGroupArgs,
21 SetRunnerSettingsArgs, Wanted,
22};
23use g1t_contracts::time::rfc3339;
24use g1t_contracts::{FailureCode, Outcome, PrincipalKind, Role, User, new_id};
25use g1t_kit::now_ms;
26use g1t_secrets::{random_hex, sha256_hex};
27use serde::Deserialize;
28use serde_json::{Map, Value, json};
29use worker::Result;
30use worker::wasm_bindgen::JsValue;
31
32use crate::plan::JobRow;
33use crate::{Actions, Count, SILENT_MS, check, fail, optional};
34
35/// How often a poll looks again for work while it waits.
36const POLL_EVERY_MS: u64 = 2_500;
37/// A runner that has been offline this long is removed, as on GitHub.
38const FORGET_OFFLINE_MS: u64 = 14 * 24 * 60 * 60 * 1000;
39/// How often a poll writes down that a runner, or the work it holds, is
40/// still there. Polls come every twenty seconds or so; a write each time
41/// would be most of the service's writes for nothing. Well inside
42/// [`model::ONLINE_WITHIN_MS`] and [`crate::SILENT_MS`].
43const SEEN_RESOLUTION_MS: u64 = 60 * 1000;
44
45/// Whether `seen` is old enough, at `at`, to be written down again.
46fn seen_is_stale(seen: Option<&str>, at: u64) -> bool {
47 seen.is_none_or(|seen| seen < rfc3339(at.saturating_sub(SEEN_RESOLUTION_MS)).as_str())
48}
49
50#[derive(Clone, Debug, Deserialize)]
51pub struct RunnerRow {
52 pub id: String,
53 pub workspace: String,
54 pub repo_id: Option<String>,
55 pub repo: Option<String>,
56 pub group_id: Option<String>,
57 pub name: String,
58 pub labels: String,
59 pub os: String,
60 pub arch: String,
61 pub version: String,
62 pub ephemeral: u32,
63 pub credential_hash: String,
64 pub previous_hash: Option<String>,
65 pub rotated_at: String,
66 pub work_id: Option<String>,
67 pub work_kind: Option<String>,
68 pub spent: u32,
69 pub last_seen_at: Option<String>,
70 pub created_at: String,
71 pub created_by: Option<String>,
72}
73
74impl RunnerRow {
75 fn labels(&self) -> Vec<String> {
76 serde_json::from_str(&self.labels).unwrap_or_default()
77 }
78
79 fn online(&self, now: u64) -> bool {
80 self.last_seen_at.as_deref().is_some_and(|seen| seen >= rfc3339(now.saturating_sub(model::ONLINE_WITHIN_MS)).as_str())
81 }
82}
83
84#[derive(Clone, Debug, Deserialize)]
85struct GroupRow {
86 id: String,
87 name: String,
88 is_default: u32,
89 repositories: String,
90 updated_at: String,
91}
92
93impl GroupRow {
94 fn repositories(&self) -> Vec<String> {
95 serde_json::from_str(&self.repositories).unwrap_or_default()
96 }
97
98 /// Whether `repo` (`owner/name`, or a bare name) may use the group.
99 fn allows(&self, repo: &str) -> bool {
100 let name = repo.rsplit('/').next().unwrap_or(repo);
101 let listed = self.repositories();
102 listed.is_empty() || listed.iter().any(|r| r.eq_ignore_ascii_case(name))
103 }
104}
105
106#[derive(Clone, Debug, Deserialize)]
107struct RegistrationRow {
108 workspace: String,
109 repo_id: Option<String>,
110 repo: Option<String>,
111 group_id: Option<String>,
112 expires_at: String,
113}
114
115#[derive(Clone, Debug, Deserialize)]
116struct SettingsRow {
117 agents: u32,
118 agent_labels: String,
119 fork_pulls: u32,
120}
121
122#[derive(Clone, Debug, Deserialize)]
123pub struct TaskRow {
124 pub id: String,
125 pub sandbox: String,
126 pub repo_id: Option<String>,
127 pub repo: String,
128 pub kind: String,
129 pub title: String,
130 pub labels: String,
131 pub env: Option<String>,
132 pub timeout_minutes: u32,
133 pub status: String,
134 pub runner_id: Option<String>,
135 pub runner_name: Option<String>,
136 pub started_at: Option<String>,
137 #[serde(default)]
138 pub seen_at: Option<String>,
139}
140
141/// Who may do what with a set of runners: the workspace, and the
142/// repository when they are a repository's own.
143struct Place {
144 workspace: String,
145 repo: Option<Repo>,
146}
147
148fn now() -> String {
149 rfc3339(now_ms())
150}
151
152fn valid_name(name: &str) -> std::result::Result<String, String> {
153 let name = name.trim();
154 let ok = !name.is_empty()
155 && name.len() <= 64
156 && name.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.'));
157 if ok {
158 Ok(name.to_owned())
159 } else {
160 Err("A runner's name is 1 to 64 letters, digits, `-`, `_` and `.`.".to_owned())
161 }
162}
163
164fn valid_group_name(name: &str) -> std::result::Result<String, String> {
165 let name = name.trim();
166 if name.is_empty() || name.len() > 64 || name.chars().any(|c| c.is_control()) {
167 return Err("A group's name is 1 to 64 characters.".to_owned());
168 }
169 Ok(name.to_owned())
170}
171
172/// What a runner's credential is checked against, and what a registration
173/// token is found by.
174fn hash(secret: &str) -> String {
175 sha256_hex(secret)
176}
177
178/// The actor of a runner's own entries in the audit log.
179fn runner_actor(runner: &RunnerRow) -> AuditActor {
180 AuditActor {
181 actor_kind: Some(g1t_contracts::audit::ActorKind::Runner),
182 actor: runner.name.clone(),
183 actor_id: runner.id.clone(),
184 ..AuditActor::default()
185 }
186}
187
188fn entry(actor: AuditActor, action: &str, workspace: &str, repo: Option<String>, path: Option<String>, message: String) -> NewAuditEntry {
189 NewAuditEntry {
190 actor,
191 action: action.to_owned(),
192 surface: Surface::Rest,
193 target: AuditTarget {
194 workspace: workspace.to_owned(),
195 repo,
196 path,
197 ..AuditTarget::default()
198 },
199 outcome: AuditOutcome::Allowed,
200 rule: "runner".to_owned(),
201 result: Some("ok".to_owned()),
202 message: Some(message),
203 request_id: new_id("req", now_ms()),
204 }
205}
206
207impl Actions {
208 // --- Shared ------------------------------------------------------------
209
210 async fn audit(&self, entries: Vec<NewAuditEntry>) {
211 let recorded: Result<u32> = g1t_kit::call(&self.events, "audit_record", &RecordAuditArgs { entries }).await;
212 if let Err(error) = recorded {
213 worker::console_error!("actions: runner audit entries not recorded: {error}");
214 }
215 }
216
217 /// Where `owner` points, if `actor` may see it (`manage` false) or
218 /// change it (`manage` true): a workspace's runners are seen by its
219 /// members and changed by its owners; a repository's own, by its
220 /// admins. Agents and workspace tokens never change them, so a
221 /// workflow's `G1T_TOKEN` cannot add a machine to run its own jobs.
222 async fn runner_place(&self, actor: &User, owner: &RunnersOwner, manage: bool) -> Result<Outcome<Place>> {
223 if actor.kind == PrincipalKind::Agent {
224 return Ok(fail(FailureCode::Forbidden, "Agents cannot see or change self-hosted runners."));
225 }
226 if manage && actor.kind == PrincipalKind::Workspace {
227 return Ok(fail(
228 FailureCode::Forbidden,
229 "A workspace's tokens, G1T_TOKEN included, cannot change self-hosted runners. Use a person's token with the runners:admin scope.",
230 ));
231 }
232 match (&owner.repo, &owner.workspace) {
233 (Some(path), _) => {
234 let repo = check!(self.may(actor, path, Capability::ManageIntegrations).await?);
235 Ok(Outcome::Ok(Place { workspace: repo.namespace.to_lowercase(), repo: Some(repo) }))
236 }
237 (None, Some(slug)) => {
238 let slug = slug.to_lowercase();
239 let role = actor.workspaces.iter().find(|m| m.slug.eq_ignore_ascii_case(&slug)).map(|m| m.role);
240 match role {
241 None => Ok(fail(FailureCode::NotFound, "There is no such workspace, or you are not a member of it.")),
242 // A workspace's machines and their tokens are its owners' alone, to see as well as change.
243 Some(Role::Member) => Ok(fail(
244 FailureCode::Forbidden,
245 format!("Only owners of {slug} can {} its self-hosted runners.", if manage { "change" } else { "see" }),
246 )),
247 Some(_) => Ok(Outcome::Ok(Place { workspace: slug, repo: None })),
248 }
249 }
250 (None, None) => Ok(fail(FailureCode::Invalid, "Give `repo` or `workspace`.")),
251 }
252 }
253
254 async fn group_rows(&self, workspace: &str) -> Result<Vec<GroupRow>> {
255 self.db
256 .prepare("SELECT * FROM runner_groups WHERE workspace = ? ORDER BY is_default DESC, name")
257 .bind(&[workspace.into()])?
258 .all()
259 .await?
260 .results::<GroupRow>()
261 }
262
263 /// The workspace's default group, made the first time it is needed.
264 async fn default_group(&self, workspace: &str) -> Result<GroupRow> {
265 let at = now();
266 self.db
267 .prepare(
268 "INSERT OR IGNORE INTO runner_groups (id, workspace, name, is_default, repositories, created_at, updated_at)
269 VALUES (?, ?, 'Default', 1, '[]', ?, ?)",
270 )
271 .bind(&[new_id("rng", now_ms()).into(), workspace.into(), at.as_str().into(), at.as_str().into()])?
272 .run()
273 .await?;
274 let row = self
275 .db
276 .prepare("SELECT * FROM runner_groups WHERE workspace = ? AND is_default = 1 LIMIT 1")
277 .bind(&[workspace.into()])?
278 .first::<GroupRow>(None)
279 .await?;
280 row.ok_or_else(|| worker::Error::RustError("the default runner group could not be made".into()))
281 }
282
283 /// A group of the workspace by id or name.
284 async fn find_group(&self, workspace: &str, wanted: &str) -> Result<Option<GroupRow>> {
285 self.default_group(workspace).await?;
286 Ok(self
287 .group_rows(workspace)
288 .await?
289 .into_iter()
290 .find(|g| g.id == wanted || g.name.eq_ignore_ascii_case(wanted.trim())))
291 }
292
293 async fn runner_rows(&self, workspace: &str) -> Result<Vec<RunnerRow>> {
294 self.db
295 .prepare("SELECT * FROM runners WHERE workspace = ? ORDER BY name")
296 .bind(&[workspace.into()])?
297 .all()
298 .await?
299 .results::<RunnerRow>()
300 }
301
302 async fn runner_by_id(&self, id: &str) -> Result<Option<RunnerRow>> {
303 self.db.prepare("SELECT * FROM runners WHERE id = ?").bind(&[id.into()])?.first::<RunnerRow>(None).await
304 }
305
306 /// What a runner is doing, for people.
307 async fn work_of(&self, row: &RunnerRow) -> Result<Option<RunnerWork>> {
308 let Some(id) = &row.work_id else { return Ok(None) };
309 if row.work_kind.as_deref() == Some("agent") {
310 let task = self.db.prepare("SELECT * FROM runner_tasks WHERE id = ?").bind(&[id.as_str().into()])?.first::<TaskRow>(None).await?;
311 return Ok(task.filter(|t| t.status == "in_progress").map(|t| RunnerWork {
312 kind: "agent".into(),
313 id: t.id,
314 name: t.title,
315 repo: Some(t.repo),
316 run_id: None,
317 started_at: t.started_at,
318 }));
319 }
320 #[derive(Deserialize)]
321 struct Found {
322 id: String,
323 name: String,
324 run_id: String,
325 status: String,
326 started_at: Option<String>,
327 repo: String,
328 }
329 let job = self
330 .db
331 .prepare("SELECT jobs.id, jobs.name, jobs.run_id, jobs.status, jobs.started_at, runs.repo FROM jobs JOIN runs ON runs.id = jobs.run_id WHERE jobs.id = ?")
332 .bind(&[id.as_str().into()])?
333 .first::<Found>(None)
334 .await?;
335 Ok(job.filter(|j| j.status == "in_progress").map(|j| RunnerWork {
336 kind: "workflow".into(),
337 id: j.id,
338 name: j.name,
339 repo: Some(j.repo),
340 run_id: Some(j.run_id),
341 started_at: j.started_at,
342 }))
343 }
344
345 async fn runner_view(&self, row: &RunnerRow, groups: &[GroupRow], at: u64) -> Result<Runner> {
346 let online = row.online(at);
347 let work = if online { self.work_of(row).await? } else { None };
348 Ok(Runner {
349 id: row.id.clone(),
350 name: row.name.clone(),
351 workspace: row.workspace.clone(),
352 repo: row.repo.clone(),
353 group: row
354 .group_id
355 .as_ref()
356 .and_then(|id| groups.iter().find(|g| &g.id == id))
357 .map(|g| g.name.clone()),
358 labels: row.labels(),
359 os: row.os.clone(),
360 arch: row.arch.clone(),
361 version: row.version.clone(),
362 ephemeral: row.ephemeral != 0,
363 status: if !online {
364 "offline"
365 } else if work.is_some() {
366 "busy"
367 } else {
368 "online"
369 }
370 .to_owned(),
371 work,
372 last_seen_at: row.last_seen_at.clone(),
373 created_at: row.created_at.clone(),
374 created_by: row.created_by.clone(),
375 })
376 }
377
378 // --- People's side -----------------------------------------------------
379
380 /// `runners`: a workspace's runners, or a repository's: its own and
381 /// the workspace's that its group lets it use.
382 pub async fn runners(&self, a: ListRunnersArgs) -> Result<Outcome<Vec<Runner>>> {
383 let place = check!(self.runner_place(&a.actor, &a.owner, false).await?);
384 self.default_group(&place.workspace).await?;
385 let groups = self.group_rows(&place.workspace).await?;
386 let at = now_ms();
387 let mut out = Vec::new();
388 for row in self.runner_rows(&place.workspace).await? {
389 let shown = match &place.repo {
390 None => row.repo_id.is_none(),
391 Some(repo) => match &row.repo_id {
392 Some(id) => id == &repo.id,
393 None => groups.iter().find(|g| Some(&g.id) == row.group_id.as_ref()).is_none_or(|g| g.allows(&repo.name)),
394 },
395 };
396 if shown {
397 out.push(self.runner_view(&row, &groups, at).await?);
398 }
399 }
400 Ok(Outcome::Ok(out))
401 }
402
403 /// `create_registration_token`.
404 pub async fn create_registration_token(&self, a: CreateRegistrationTokenArgs) -> Result<Outcome<RegistrationToken>> {
405 let place = check!(self.runner_place(&a.actor, &a.owner, true).await?);
406 let made = self
407 .db
408 .prepare("SELECT COUNT(*) AS n FROM runner_registrations WHERE workspace = ? AND created_at > ?")
409 .bind(&[place.workspace.as_str().into(), rfc3339(now_ms().saturating_sub(60 * 60 * 1000)).into()])?
410 .first::<Count>(None)
411 .await?
412 .map_or(0, |c| c.n);
413 if made >= model::MAX_TOKENS_PER_HOUR {
414 return Ok(fail(FailureCode::Conflict, "Too many registration tokens in the last hour. Each lasts an hour and registers any number of runners."));
415 }
416 let group = match (&place.repo, a.group.as_deref().map(str::trim).filter(|g| !g.is_empty())) {
417 (Some(_), Some(_)) => return Ok(fail(FailureCode::Invalid, "A repository's own runners are in no group.")),
418 (Some(_), None) => None,
419 (None, Some(wanted)) => match self.find_group(&place.workspace, wanted).await? {
420 Some(group) => Some(group),
421 None => return Ok(fail(FailureCode::NotFound, format!("There is no runner group called {wanted}."))),
422 },
423 (None, None) => Some(self.default_group(&place.workspace).await?),
424 };
425 let token = format!("{}{}", model::REGISTRATION_PREFIX, random_hex(24));
426 let at = now_ms();
427 let expires_at = rfc3339(at + model::REGISTRATION_TTL_SECONDS * 1000);
428 let repo_name = place.repo.as_ref().map(|r| format!("{}/{}", r.namespace, r.name));
429 self.db
430 .prepare(
431 "INSERT INTO runner_registrations (token_hash, workspace, repo_id, repo, group_id, created_by, created_at, expires_at)
432 VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
433 )
434 .bind(&[
435 hash(&token).into(),
436 place.workspace.as_str().into(),
437 optional(place.repo.as_ref().map(|r| r.id.as_str())),
438 optional(repo_name.as_deref()),
439 optional(group.as_ref().map(|g| g.id.as_str())),
440 a.actor.username.as_str().into(),
441 rfc3339(at).into(),
442 expires_at.as_str().into(),
443 ])?
444 .run()
445 .await?;
446 self.audit(vec![entry(
447 AuditActor::of(&a.actor),
448 "create_runner_registration_token",
449 &place.workspace,
450 repo_name.clone(),
451 None,
452 format!("Made a runner registration token{}.", group.as_ref().map(|g| format!(" for group {}", g.name)).unwrap_or_default()),
453 )])
454 .await;
455 Ok(Outcome::Ok(RegistrationToken {
456 token,
457 expires_at,
458 workspace: place.workspace,
459 repo: repo_name,
460 group: group.map(|g| g.name),
461 url: crate::SITE.to_owned(),
462 }))
463 }
464
465 /// `remove_runner`: a person removing one.
466 pub async fn remove_runner(&self, a: RemoveRunnerArgs) -> Result<Outcome<bool>> {
467 let place = check!(self.runner_place(&a.actor, &a.owner, true).await?);
468 let Some(row) = self.runner_by_id(&a.id).await? else {
469 return Ok(fail(FailureCode::NotFound, "There is no such runner."));
470 };
471 let ours = row.workspace == place.workspace
472 && match &place.repo {
473 Some(repo) => row.repo_id.as_deref() == Some(repo.id.as_str()),
474 None => row.repo_id.is_none(),
475 };
476 if !ours {
477 return Ok(fail(FailureCode::NotFound, "There is no such runner."));
478 }
479 self.forget_runner(&row, &format!("{} removed the runner {}.", a.actor.username, row.name)).await?;
480 self.audit(vec![entry(
481 AuditActor::of(&a.actor),
482 "remove_runner",
483 &row.workspace,
484 row.repo.clone(),
485 Some(row.name.clone()),
486 format!("Removed the self-hosted runner {}.", row.name),
487 )])
488 .await;
489 Ok(Outcome::Ok(true))
490 }
491
492 /// Removes a runner: what it was running fails, and its credential
493 /// stops working at once.
494 async fn forget_runner(&self, row: &RunnerRow, why: &str) -> Result<()> {
495 self.db.prepare("DELETE FROM runners WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
496 self.abandon_work(&row.id, why).await
497 }
498
499 /// Fails whatever a runner was running.
500 async fn abandon_work(&self, runner_id: &str, why: &str) -> Result<()> {
501 let jobs = self
502 .db
503 .prepare("SELECT id FROM jobs WHERE runner_id = ? AND status = 'in_progress'")
504 .bind(&[runner_id.into()])?
505 .all()
506 .await?
507 .results::<IdRow>()?;
508 for job in jobs {
509 self.finish_job(&job.id, "failure", Some(why), None).await?;
510 }
511 let tasks = self
512 .db
513 .prepare("SELECT * FROM runner_tasks WHERE runner_id = ? AND status = 'in_progress'")
514 .bind(&[runner_id.into()])?
515 .all()
516 .await?
517 .results::<TaskRow>()?;
518 for task in tasks {
519 self.end_task(&task, 1, why).await?;
520 }
521 Ok(())
522 }
523
524 /// `runner_groups`.
525 pub async fn runner_groups(&self, a: RunnerGroupsArgs) -> Result<Outcome<Vec<RunnerGroup>>> {
526 let owner = RunnersOwner { repo: None, workspace: Some(a.workspace.clone()) };
527 let place = check!(self.runner_place(&a.actor, &owner, false).await?);
528 self.default_group(&place.workspace).await?;
529 let groups = self.group_rows(&place.workspace).await?;
530 let runners = self.runner_rows(&place.workspace).await?;
531 Ok(Outcome::Ok(groups.iter().map(|g| group_view(g, &runners)).collect()))
532 }
533
534 /// `set_runner_group`.
535 pub async fn set_runner_group(&self, a: SetRunnerGroupArgs) -> Result<Outcome<RunnerGroup>> {
536 let owner = RunnersOwner { repo: None, workspace: Some(a.workspace.clone()) };
537 let place = check!(self.runner_place(&a.actor, &owner, true).await?);
538 self.default_group(&place.workspace).await?;
539 let groups = self.group_rows(&place.workspace).await?;
540 let name = match a.name.as_deref().map(valid_group_name).transpose() {
541 Ok(name) => name,
542 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
543 };
544 let repositories: Option<Vec<String>> = a.repositories.map(|list| {
545 let mut out: Vec<String> = Vec::new();
546 for repo in list {
547 let name = repo.trim().rsplit('/').next().unwrap_or_default().to_owned();
548 if !name.is_empty() && !out.iter().any(|r: &String| r.eq_ignore_ascii_case(&name)) {
549 out.push(name);
550 }
551 }
552 out
553 });
554 if let Some(name) = &name
555 && groups.iter().any(|g| g.name.eq_ignore_ascii_case(name) && Some(&g.id) != a.id.as_ref())
556 {
557 return Ok(fail(FailureCode::Conflict, format!("There is already a group called {name}.")));
558 }
559 let at = now();
560 let id = match &a.id {
561 Some(id) => {
562 let Some(group) = groups.iter().find(|g| &g.id == id) else {
563 return Ok(fail(FailureCode::NotFound, "There is no such runner group."));
564 };
565 self.db
566 .prepare("UPDATE runner_groups SET name = ?, repositories = ?, updated_at = ? WHERE id = ?")
567 .bind(&[
568 name.clone().unwrap_or(group.name.clone()).into(),
569 serde_json::to_string(&repositories.clone().unwrap_or(group.repositories()))?.into(),
570 at.as_str().into(),
571 id.as_str().into(),
572 ])?
573 .run()
574 .await?;
575 id.clone()
576 }
577 None => {
578 let Some(name) = name.clone() else {
579 return Ok(fail(FailureCode::Invalid, "A new group needs a name."));
580 };
581 if groups.len() >= 50 {
582 return Ok(fail(FailureCode::Conflict, "A workspace has at most 50 runner groups."));
583 }
584 let id = new_id("rng", now_ms());
585 self.db
586 .prepare(
587 "INSERT INTO runner_groups (id, workspace, name, is_default, repositories, created_at, updated_at) VALUES (?, ?, ?, 0, ?, ?, ?)",
588 )
589 .bind(&[
590 id.as_str().into(),
591 place.workspace.as_str().into(),
592 name.into(),
593 serde_json::to_string(&repositories.clone().unwrap_or_default())?.into(),
594 at.as_str().into(),
595 at.as_str().into(),
596 ])?
597 .run()
598 .await?;
599 id
600 }
601 };
602 let groups = self.group_rows(&place.workspace).await?;
603 let runners = self.runner_rows(&place.workspace).await?;
604 let Some(group) = groups.iter().find(|g| g.id == id) else {
605 return Ok(fail(FailureCode::NotFound, "There is no such runner group."));
606 };
607 self.audit(vec![entry(
608 AuditActor::of(&a.actor),
609 if a.id.is_some() { "update_runner_group" } else { "create_runner_group" },
610 &place.workspace,
611 None,
612 Some(group.name.clone()),
613 format!(
614 "Runner group {}: {}.",
615 group.name,
616 if group.repositories().is_empty() { "every repository".to_owned() } else { group.repositories().join(", ") }
617 ),
618 )])
619 .await;
620 Ok(Outcome::Ok(group_view(group, &runners)))
621 }
622
623 /// `delete_runner_group`.
624 pub async fn delete_runner_group(&self, a: DeleteRunnerGroupArgs) -> Result<Outcome<bool>> {
625 let owner = RunnersOwner { repo: None, workspace: Some(a.workspace.clone()) };
626 let place = check!(self.runner_place(&a.actor, &owner, true).await?);
627 let default = self.default_group(&place.workspace).await?;
628 let Some(group) = self.group_rows(&place.workspace).await?.into_iter().find(|g| g.id == a.id) else {
629 return Ok(fail(FailureCode::NotFound, "There is no such runner group."));
630 };
631 if group.is_default != 0 {
632 return Ok(fail(FailureCode::Conflict, "The default group cannot be deleted."));
633 }
634 self.db
635 .batch(vec![
636 self.db.prepare("UPDATE runners SET group_id = ? WHERE group_id = ?").bind(&[default.id.as_str().into(), group.id.as_str().into()])?,
637 self.db.prepare("DELETE FROM runner_groups WHERE id = ?").bind(&[group.id.as_str().into()])?,
638 ])
639 .await?;
640 self.audit(vec![entry(
641 AuditActor::of(&a.actor),
642 "delete_runner_group",
643 &place.workspace,
644 None,
645 Some(group.name.clone()),
646 format!("Deleted the runner group {}; its runners joined {}.", group.name, default.name),
647 )])
648 .await;
649 Ok(Outcome::Ok(true))
650 }
651
652 async fn settings_row(&self, owner: &str) -> Result<Option<SettingsRow>> {
653 self.db.prepare("SELECT * FROM runner_settings WHERE owner = ?").bind(&[owner.into()])?.first::<SettingsRow>(None).await
654 }
655
656 fn settings_of(row: Option<SettingsRow>, inherited: bool) -> RunnerSettings {
657 match row {
658 Some(row) => RunnerSettings {
659 agents_on_self_hosted: row.agents != 0,
660 agent_labels: serde_json::from_str(&row.agent_labels).unwrap_or_else(|_| vec![model::SELF_HOSTED.to_owned()]),
661 fork_pull_requests: row.fork_pulls != 0,
662 inherited,
663 },
664 None => RunnerSettings {
665 agent_labels: vec![model::SELF_HOSTED.to_owned()],
666 inherited,
667 ..RunnerSettings::default()
668 },
669 }
670 }
671
672 /// The settings that apply in a workspace, or in one of its
673 /// repositories: the repository's own if it has them.
674 pub async fn effective_runner_settings(&self, workspace: &str, repo_id: Option<&str>) -> Result<RunnerSettings> {
675 if let Some(repo_id) = repo_id
676 && let Some(own) = self.settings_row(repo_id).await?
677 {
678 return Ok(Self::settings_of(Some(own), false));
679 }
680 let inherited = repo_id.is_some();
681 Ok(Self::settings_of(self.settings_row(&workspace.to_lowercase()).await?, inherited))
682 }
683
684 /// `runner_settings`.
685 pub async fn runner_settings(&self, a: RunnerSettingsArgs) -> Result<Outcome<RunnerSettings>> {
686 let place = check!(self.runner_place(&a.actor, &a.owner, false).await?);
687 Ok(Outcome::Ok(self.effective_runner_settings(&place.workspace, place.repo.as_ref().map(|r| r.id.as_str())).await?))
688 }
689
690 /// `set_runner_settings`.
691 pub async fn set_runner_settings(&self, a: SetRunnerSettingsArgs) -> Result<Outcome<RunnerSettings>> {
692 let place = check!(self.runner_place(&a.actor, &a.owner, true).await?);
693 let (owner, scope) = match &place.repo {
694 Some(repo) => (repo.id.clone(), "repository"),
695 None => (place.workspace.clone(), "workspace"),
696 };
697 if a.inherit && place.repo.is_some() {
698 self.db.prepare("DELETE FROM runner_settings WHERE owner = ?").bind(&[owner.as_str().into()])?.run().await?;
699 } else {
700 let current = self.effective_runner_settings(&place.workspace, place.repo.as_ref().map(|r| r.id.as_str())).await?;
701 let labels = match &a.agent_labels {
702 Some(given) => match model::runner_labels(given, model::SELF_HOSTED, model::SELF_HOSTED) {
703 Ok(labels) => labels,
704 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
705 },
706 None => current.agent_labels.clone(),
707 };
708 let agents = a.agents_on_self_hosted.unwrap_or(current.agents_on_self_hosted);
709 let forks = a.fork_pull_requests.unwrap_or(current.fork_pull_requests);
710 self.db
711 .prepare(
712 "INSERT INTO runner_settings (owner, scope, agents, agent_labels, fork_pulls, updated_at, updated_by)
713 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
714 ON CONFLICT (owner) DO UPDATE SET agents = ?3, agent_labels = ?4, fork_pulls = ?5, updated_at = ?6, updated_by = ?7",
715 )
716 .bind(&[
717 owner.as_str().into(),
718 scope.into(),
719 u32::from(agents).into(),
720 serde_json::to_string(&labels)?.into(),
721 u32::from(forks).into(),
722 now().into(),
723 a.actor.username.as_str().into(),
724 ])?
725 .run()
726 .await?;
727 }
728 let settings = self.effective_runner_settings(&place.workspace, place.repo.as_ref().map(|r| r.id.as_str())).await?;
729 self.audit(vec![entry(
730 AuditActor::of(&a.actor),
731 "update_runner_settings",
732 &place.workspace,
733 place.repo.as_ref().map(|r| format!("{}/{}", r.namespace, r.name)),
734 None,
735 format!(
736 "Agents on self-hosted runners: {}; labels {}; pull requests from forks: {}.",
737 if settings.agents_on_self_hosted { "on" } else { "off" },
738 settings.agent_labels.join(", "),
739 if settings.fork_pull_requests { "allowed" } else { "not allowed" }
740 ),
741 )])
742 .await;
743 Ok(Outcome::Ok(settings))
744 }
745
746 // --- The runner's side ---------------------------------------------------
747
748 /// `runner_register`.
749 pub async fn runner_register(&self, a: RegisterArgs) -> Result<Outcome<Registered>> {
750 let refused = || fail(FailureCode::Unauthenticated, "That registration token is not valid, or it expired. Make a new one under Settings, Runners.");
751 if !a.token.starts_with(model::REGISTRATION_PREFIX) {
752 return Ok(refused());
753 }
754 let registration = self
755 .db
756 .prepare("SELECT * FROM runner_registrations WHERE token_hash = ?")
757 .bind(&[hash(&a.token).into()])?
758 .first::<RegistrationRow>(None)
759 .await?;
760 let at = now_ms();
761 let Some(registration) = registration.filter(|r| r.expires_at.as_str() > rfc3339(at).as_str()) else {
762 return Ok(refused());
763 };
764 let name = match valid_name(&a.name) {
765 Ok(name) => name,
766 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
767 };
768 let (Some(os), Some(arch)) = (model::os_label(&a.os), model::arch_label(&a.arch)) else {
769 return Ok(fail(FailureCode::Invalid, "A runner's os is linux, macos or windows, and its arch x64 or arm64."));
770 };
771 let labels = match model::runner_labels(&a.labels, os, arch) {
772 Ok(labels) => labels,
773 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
774 };
775 let workspace = registration.workspace.clone();
776 // Abuse: so many machines, so fast, are not a team's.
777 let recent = self
778 .db
779 .prepare("SELECT COUNT(*) AS n FROM runners WHERE workspace = ? AND created_at > ?")
780 .bind(&[workspace.as_str().into(), rfc3339(at.saturating_sub(60 * 1000)).into()])?
781 .first::<Count>(None)
782 .await?
783 .map_or(0, |c| c.n);
784 if recent >= model::MAX_REGISTRATIONS_PER_MINUTE {
785 return Ok(fail(FailureCode::Conflict, "Too many runners registered in the last minute. Wait a moment and try again."));
786 }
787 let total = self
788 .db
789 .prepare("SELECT COUNT(*) AS n FROM runners WHERE workspace = ?")
790 .bind(&[workspace.as_str().into()])?
791 .first::<Count>(None)
792 .await?
793 .map_or(0, |c| c.n);
794 // A group: the one asked for, the token's, or the default.
795 let group = if registration.repo_id.is_some() {
796 None
797 } else {
798 match a.group.as_deref().map(str::trim).filter(|g| !g.is_empty()) {
799 Some(wanted) => match self.find_group(&workspace, wanted).await? {
800 Some(group) => Some(group.id),
801 None => return Ok(fail(FailureCode::NotFound, format!("There is no runner group called {wanted}."))),
802 },
803 None => match registration.group_id.clone() {
804 Some(id) => Some(id),
805 None => Some(self.default_group(&workspace).await?.id),
806 },
807 }
808 };
809 let existing = self
810 .db
811 .prepare("SELECT * FROM runners WHERE workspace = ? AND COALESCE(repo_id, '') = ? AND lower(name) = lower(?)")
812 .bind(&[workspace.as_str().into(), registration.repo_id.clone().unwrap_or_default().into(), name.as_str().into()])?
813 .first::<RunnerRow>(None)
814 .await?;
815 if let Some(existing) = &existing {
816 if !a.replace {
817 return Ok(fail(
818 FailureCode::Conflict,
819 format!("A runner called {name} is already registered here. Register with --replace to take its place, or choose another --name."),
820 ));
821 }
822 self.forget_runner(existing, &format!("The runner {name} was replaced by a new registration.")).await?;
823 } else if total >= model::MAX_RUNNERS {
824 return Ok(fail(FailureCode::Conflict, format!("A workspace has at most {} self-hosted runners.", model::MAX_RUNNERS)));
825 }
826 let id = new_id("rnr", at);
827 let credential = format!("{}{}", model::CREDENTIAL_PREFIX, random_hex(32));
828 let created = rfc3339(at);
829 self.db
830 .prepare(
831 "INSERT INTO runners (id, workspace, repo_id, repo, group_id, name, labels, os, arch, version, ephemeral, credential_hash,
832 rotated_at, last_seen_at, created_at, created_by)
833 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
834 )
835 .bind(&[
836 id.as_str().into(),
837 workspace.as_str().into(),
838 optional(registration.repo_id.as_deref()),
839 optional(registration.repo.as_deref()),
840 optional(group.as_deref()),
841 name.as_str().into(),
842 serde_json::to_string(&labels)?.into(),
843 os.into(),
844 arch.into(),
845 a.version.chars().take(40).collect::<String>().into(),
846 u32::from(a.ephemeral).into(),
847 hash(&credential).into(),
848 created.as_str().into(),
849 created.as_str().into(),
850 created.as_str().into(),
851 JsValue::NULL,
852 ])?
853 .run()
854 .await?;
855 self.db
856 .prepare("UPDATE runner_registrations SET used = used + 1 WHERE token_hash = ?")
857 .bind(&[hash(&a.token).into()])?
858 .run()
859 .await?;
860 let row = self.runner_by_id(&id).await?.ok_or_else(|| worker::Error::RustError("the runner was not kept".into()))?;
861 self.audit(vec![entry(
862 runner_actor(&row),
863 "register_runner",
864 &workspace,
865 row.repo.clone(),
866 Some(name.clone()),
867 format!("Registered the self-hosted runner {name} ({os}, {arch}; labels {}).", labels.join(", ")),
868 )])
869 .await;
870 let groups = self.group_rows(&workspace).await?;
871 Ok(Outcome::Ok(Registered {
872 runner: self.runner_view(&row, &groups, at).await?,
873 credential,
874 }))
875 }
876
877 /// The runner a call is from, if its credential is that runner's (or
878 /// the one before a rotation, until the new one is used).
879 async fn authenticated(&self, auth: &RunnerAuth) -> Result<Option<(RunnerRow, bool)>> {
880 if !auth.credential.starts_with(model::CREDENTIAL_PREFIX) {
881 return Ok(None);
882 }
883 let Some(row) = self.runner_by_id(&auth.runner).await? else { return Ok(None) };
884 let given = hash(&auth.credential);
885 if g1t_secrets::same(&row.credential_hash, &given) {
886 if row.previous_hash.is_some() {
887 self.db.prepare("UPDATE runners SET previous_hash = NULL WHERE id = ?").bind(&[row.id.as_str().into()])?.run().await?;
888 }
889 return Ok(Some((row, true)));
890 }
891 if row.previous_hash.as_deref().is_some_and(|previous| g1t_secrets::same(previous, &given)) {
892 return Ok(Some((row, false)));
893 }
894 Ok(None)
895 }
896
897 fn unknown_runner<T>() -> Outcome<T> {
898 fail(FailureCode::Unauthenticated, "This runner is not registered, or its credential is not valid. Register it again.")
899 }
900
901 /// `runner_poll`.
902 pub async fn runner_poll(&self, a: PollArgs) -> Result<Outcome<Poll>> {
903 let Some((row, current)) = self.authenticated(&a.auth).await? else {
904 return Ok(Self::unknown_runner());
905 };
906 let at = now_ms();
907 let seen = rfc3339(at);
908 // A new credential every day, handed over on a poll that used the
909 // current one.
910 let mut credential = None;
911 let rotate_after = rfc3339(at.saturating_sub(model::ROTATE_AFTER_MS));
912 if current && row.rotated_at.as_str() < rotate_after.as_str() {
913 let fresh = format!("{}{}", model::CREDENTIAL_PREFIX, random_hex(32));
914 self.db
915 .prepare("UPDATE runners SET previous_hash = credential_hash, credential_hash = ?, rotated_at = ? WHERE id = ?")
916 .bind(&[hash(&fresh).into(), seen.as_str().into(), row.id.as_str().into()])?
917 .run()
918 .await?;
919 credential = Some(fresh);
920 }
921 // What it holds that it should stop, and the heartbeat for the rest.
922 let mut cancel = Vec::new();
923 let mut active: Option<(String, &'static str)> = None;
924 for id in a.running.iter().take(10) {
925 if let Some(job) = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[id.as_str().into()])?.first::<JobRow>(None).await? {
926 if job.status == "in_progress" && job.runner_id.as_deref() == Some(row.id.as_str()) {
927 if seen_is_stale(job.seen_at.as_deref(), at) {
928 self.db.prepare("UPDATE jobs SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?;
929 }
930 active = Some((id.clone(), "workflow"));
931 } else {
932 cancel.push(id.clone());
933 }
934 continue;
935 }
936 match self.db.prepare("SELECT * FROM runner_tasks WHERE id = ?").bind(&[id.as_str().into()])?.first::<TaskRow>(None).await? {
937 Some(task) if task.status == "in_progress" && task.runner_id.as_deref() == Some(row.id.as_str()) => {
938 if seen_is_stale(task.seen_at.as_deref(), at) {
939 self.db.prepare("UPDATE runner_tasks SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?;
940 }
941 active = Some((id.clone(), "agent"));
942 }
943 _ => cancel.push(id.clone()),
944 }
945 }
946 // Written when something changed, or once a minute to keep it online.
947 let version: String = if a.version.is_empty() { row.version.clone() } else { a.version.chars().take(40).collect() };
948 let work_id = active.as_ref().map(|(id, _)| id.as_str());
949 let work_kind = active.as_ref().map(|(_, kind)| *kind);
950 let changed = version != row.version || work_id != row.work_id.as_deref() || work_kind != row.work_kind.as_deref();
951 if changed || seen_is_stale(row.last_seen_at.as_deref(), at) {
952 self.db
953 .prepare("UPDATE runners SET last_seen_at = ?, version = ?, work_id = ?, work_kind = ? WHERE id = ?")
954 .bind(&[seen.as_str().into(), version.into(), optional(work_id), optional(work_kind), row.id.as_str().into()])?
955 .run()
956 .await?;
957 }
958 let mut poll = Poll { assignment: None, cancel, credential, removed: false };
959 // Busy, or an ephemeral runner that has had its one job.
960 if active.is_some() || !a.running.is_empty() {
961 return Ok(Outcome::Ok(poll));
962 }
963 if row.ephemeral != 0 && row.spent != 0 {
964 self.forget_runner(&row, "The ephemeral runner finished its job.").await?;
965 poll.removed = true;
966 return Ok(Outcome::Ok(poll));
967 }
968 let groups = self.group_rows(&row.workspace).await?;
969 let group = row.group_id.as_ref().and_then(|id| groups.iter().find(|g| &g.id == id)).cloned();
970 let deadline = at + a.wait_ms.min(model::MAX_POLL_WAIT_MS);
971 loop {
972 if let Some(assignment) = self.claim(&row, group.as_ref()).await? {
973 poll.assignment = Some(assignment);
974 return Ok(Outcome::Ok(poll));
975 }
976 if now_ms() + POLL_EVERY_MS > deadline {
977 return Ok(Outcome::Ok(poll));
978 }
979 worker::Delay::from(std::time::Duration::from_millis(POLL_EVERY_MS)).await;
980 }
981 }
982
983 /// Takes the oldest work this runner can do, if there is any.
984 async fn claim(&self, runner: &RunnerRow, group: Option<&GroupRow>) -> Result<Option<Assignment>> {
985 let labels = runner.labels();
986 let group_name = group.map(|g| g.name.clone());
987 // Workflow jobs.
988 let mut sql = "SELECT jobs.*, runs.repo AS run_repo FROM jobs JOIN runs ON runs.id = jobs.run_id
989 WHERE jobs.status = 'queued' AND jobs.labels IS NOT NULL AND lower(jobs.namespace) = ?"
990 .to_owned();
991 let mut binds: Vec<JsValue> = vec![runner.workspace.as_str().into()];
992 if let Some(repo_id) = &runner.repo_id {
993 sql.push_str(" AND jobs.repo_id = ?");
994 binds.push(repo_id.as_str().into());
995 }
996 // One job of a concurrency group runs at a time.
997 sql.push_str(
998 " AND (jobs.concurrency_group IS NULL OR NOT EXISTS (SELECT 1 FROM jobs o WHERE o.repo_id = jobs.repo_id
999 AND o.concurrency_group = jobs.concurrency_group AND o.status = 'in_progress'))",
1000 );
1001 sql.push_str(" ORDER BY jobs.rowid LIMIT 50");
1002 #[derive(Deserialize)]
1003 struct Queued {
1004 #[serde(flatten)]
1005 job: JobRow,
1006 run_repo: String,
1007 }
1008 let queued = self.db.prepare(sql).bind(&binds)?.all().await?.results::<Queued>()?;
1009 for Queued { job, run_repo } in queued {
1010 let stored: Vec<String> = job.labels.as_deref().and_then(|l| serde_json::from_str(l).ok()).unwrap_or_default();
1011 let wanted = Wanted::from_stored(&stored);
1012 if !wanted.matches(&labels, group_name.as_deref()) {
1013 continue;
1014 }
1015 if runner.repo_id.is_none() && !group.is_none_or(|g| g.allows(&run_repo)) {
1016 continue;
1017 }
1018 if let Some(max) = job.max_parallel {
1019 let siblings = self
1020 .db
1021 .prepare("SELECT COUNT(*) AS n FROM jobs WHERE run_id = ? AND key = ? AND status = 'in_progress'")
1022 .bind(&[job.run_id.as_str().into(), job.key.as_str().into()])?
1023 .first::<Count>(None)
1024 .await?
1025 .map_or(0, |c| c.n);
1026 if siblings >= max {
1027 continue;
1028 }
1029 }
1030 let token = random_hex(24);
1031 let at = now();
1032 let claimed = self
1033 .db
1034 .prepare(
1035 "UPDATE jobs SET status = 'in_progress', token_hash = ?, runner_id = ?, runner_name = ?, reason = NULL, started_at = ?, seen_at = ?
1036 WHERE id = ? AND status = 'queued' RETURNING id",
1037 )
1038 .bind(&[
1039 sha256_hex(&token).into(),
1040 runner.id.as_str().into(),
1041 runner.name.as_str().into(),
1042 at.as_str().into(),
1043 at.as_str().into(),
1044 job.id.as_str().into(),
1045 ])?
1046 .first::<Value>(None)
1047 .await?;
1048 if claimed.is_none() {
1049 continue;
1050 }
1051 self.db
1052 .batch(vec![
1053 self.db
1054 .prepare("UPDATE runs SET status = 'in_progress', started_at = COALESCE(started_at, ?) WHERE id = ? AND status = 'queued'")
1055 .bind(&[at.as_str().into(), job.run_id.as_str().into()])?,
1056 self.db
1057 .prepare("UPDATE runners SET work_id = ?, work_kind = 'workflow', spent = ephemeral WHERE id = ?")
1058 .bind(&[job.id.as_str().into(), runner.id.as_str().into()])?,
1059 ])
1060 .await?;
1061 self.audit(vec![entry(
1062 runner_actor(runner),
1063 "runner.take_job",
1064 &runner.workspace,
1065 Some(run_repo.clone()),
1066 Some(job.name.clone()),
1067 format!("The self-hosted runner {} took the job {}.", runner.name, job.name),
1068 )])
1069 .await;
1070 // A job that deploys: its run's deployment is under way.
1071 self.job_started(&job.id).await?;
1072 let image = self.job_image(&job).await?;
1073 return Ok(Some(Assignment {
1074 kind: "workflow".into(),
1075 id: job.id.clone(),
1076 name: job.name.clone(),
1077 repo: run_repo,
1078 timeout_minutes: job.timeout_minutes,
1079 image,
1080 token: Some(token),
1081 env: None,
1082 }));
1083 }
1084
1085 // Agent work.
1086 let tasks = self
1087 .db
1088 .prepare("SELECT * FROM runner_tasks WHERE status = 'queued' AND workspace = ? ORDER BY created_at LIMIT 20")
1089 .bind(&[runner.workspace.as_str().into()])?
1090 .all()
1091 .await?
1092 .results::<TaskRow>()?;
1093 for task in tasks {
1094 let wanted = Wanted::of(&serde_json::from_str::<Vec<String>>(&task.labels).unwrap_or_default(), None);
1095 if !wanted.matches(&labels, None) {
1096 continue;
1097 }
1098 if let Some(repo_id) = &runner.repo_id
1099 && task.repo_id.as_deref() != Some(repo_id.as_str())
1100 {
1101 continue;
1102 }
1103 if runner.repo_id.is_none() && !group.is_none_or(|g| g.allows(&task.repo)) {
1104 continue;
1105 }
1106 let Some(sealed) = task.env.clone() else { continue };
1107 let at = now();
1108 let claimed = self
1109 .db
1110 .prepare(
1111 "UPDATE runner_tasks SET status = 'in_progress', runner_id = ?, runner_name = ?, env = NULL, started_at = ?, seen_at = ?
1112 WHERE id = ? AND status = 'queued' RETURNING id",
1113 )
1114 .bind(&[runner.id.as_str().into(), runner.name.as_str().into(), at.as_str().into(), at.as_str().into(), task.id.as_str().into()])?
1115 .first::<Value>(None)
1116 .await?;
1117 if claimed.is_none() {
1118 continue;
1119 }
1120 self.db
1121 .prepare("UPDATE runners SET work_id = ?, work_kind = 'agent', spent = ephemeral WHERE id = ?")
1122 .bind(&[task.id.as_str().into(), runner.id.as_str().into()])?
1123 .run()
1124 .await?;
1125 let env: Map<String, Value> = self
1126 .sealer
1127 .as_ref()
1128 .and_then(|sealer| sealer.open(&sealed, &task.id))
1129 .and_then(|text| serde_json::from_str(&text).ok())
1130 .unwrap_or_default();
1131 if env.is_empty() {
1132 self.end_task(&TaskRow { status: "in_progress".into(), ..task.clone() }, 1, "The work's environment could not be opened.").await?;
1133 continue;
1134 }
1135 self.audit(vec![entry(
1136 runner_actor(runner),
1137 "runner.take_task",
1138 &runner.workspace,
1139 Some(task.repo.clone()),
1140 Some(task.kind.clone()),
1141 format!("The self-hosted runner {} took agent work: {}.", runner.name, task.title),
1142 )])
1143 .await;
1144 return Ok(Some(Assignment {
1145 kind: "agent".into(),
1146 id: task.id.clone(),
1147 name: task.title.clone(),
1148 repo: task.repo.clone(),
1149 timeout_minutes: task.timeout_minutes,
1150 image: None,
1151 token: None,
1152 env: Some(env),
1153 }));
1154 }
1155 Ok(None)
1156 }
1157
1158 /// The image a job's `container:` names, when it names one plainly.
1159 async fn job_image(&self, job: &JobRow) -> Result<Option<String>> {
1160 let Some(run) = self.run_row(&job.run_id).await? else { return Ok(None) };
1161 let raw = match job.callee() {
1162 Some((_, spec, _)) => spec.raw,
1163 None => match g1t_actions::workflow::parse(&run.source) {
1164 Ok(workflow) => match workflow.jobs.into_iter().find(|j| j.id == job.key) {
1165 Some(spec) => spec.raw,
1166 None => return Ok(None),
1167 },
1168 Err(_) => return Ok(None),
1169 },
1170 };
1171 let image = match raw.get("container") {
1172 Some(Value::String(image)) => Some(image.clone()),
1173 Some(Value::Object(container)) => container.get("image").and_then(Value::as_str).map(str::to_owned),
1174 _ => None,
1175 };
1176 Ok(image.filter(|image| !image.contains("${{") && !image.trim().is_empty()))
1177 }
1178
1179 /// `runner_finished`.
1180 pub async fn runner_finished(&self, a: FinishedArgs) -> Result<Outcome<bool>> {
1181 let Some((row, _)) = self.authenticated(&a.auth).await? else {
1182 return Ok(Self::unknown_runner());
1183 };
1184 let reason = a.reason.clone().filter(|r| !r.trim().is_empty()).map(|r| r.chars().take(2000).collect::<String>());
1185 let mut changed = false;
1186 if let Some(job) = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[a.id.as_str().into()])?.first::<JobRow>(None).await? {
1187 if job.runner_id.as_deref() == Some(row.id.as_str()) && job.status == "in_progress" {
1188 let why = reason.unwrap_or_else(|| format!("The job's process on {} exited with code {} before it reported how the job went.", row.name, a.exit_code));
1189 self.finish_job(&job.id, "failure", Some(&why), None).await?;
1190 changed = true;
1191 }
1192 } else if let Some(task) = self.db.prepare("SELECT * FROM runner_tasks WHERE id = ?").bind(&[a.id.as_str().into()])?.first::<TaskRow>(None).await?
1193 && task.runner_id.as_deref() == Some(row.id.as_str())
1194 && task.status == "in_progress"
1195 {
1196 let why = reason.unwrap_or_else(|| if a.exit_code == 0 { String::new() } else { format!("The work exited with code {} on {}.", a.exit_code, row.name) });
1197 self.end_task(&task, a.exit_code, &why).await?;
1198 changed = true;
1199 }
1200 self.db
1201 .prepare("UPDATE runners SET work_id = NULL, work_kind = NULL, last_seen_at = ? WHERE id = ? AND work_id = ?")
1202 .bind(&[now().into(), row.id.as_str().into(), a.id.as_str().into()])?
1203 .run()
1204 .await?;
1205 Ok(Outcome::Ok(changed))
1206 }
1207
1208 /// `runner_remove_self`.
1209 pub async fn runner_remove_self(&self, a: RemoveSelfArgs) -> Result<Outcome<bool>> {
1210 let Some((row, _)) = self.authenticated(&a.auth).await? else {
1211 return Ok(Self::unknown_runner());
1212 };
1213 self.forget_runner(&row, &format!("The runner {} was removed from its machine.", row.name)).await?;
1214 self.audit(vec![entry(
1215 runner_actor(&row),
1216 "remove_runner",
1217 &row.workspace,
1218 row.repo.clone(),
1219 Some(row.name.clone()),
1220 format!("The self-hosted runner {} removed itself.", row.name),
1221 )])
1222 .await;
1223 Ok(Outcome::Ok(true))
1224 }
1225
1226 /// A self-hosted job finished, however it did: its runner is free, and
1227 /// its time goes on the workspace's usage as self-hosted, at $0.
1228 pub async fn released(&self, job: &JobRow) -> Result<()> {
1229 let Some(runner_id) = &job.runner_id else { return Ok(()) };
1230 self.db
1231 .prepare("UPDATE runners SET work_id = NULL, work_kind = NULL WHERE id = ? AND work_id = ?")
1232 .bind(&[runner_id.as_str().into(), job.id.as_str().into()])?
1233 .run()
1234 .await?;
1235 let (Some(started), Some(finished)) = (job.started_at.as_deref(), job.finished_at.as_deref().map(str::to_owned).or_else(|| Some(now()))) else {
1236 return Ok(());
1237 };
1238 let seconds = match (g1t_contracts::time::parse_rfc3339(started), g1t_contracts::time::parse_rfc3339(&finished)) {
1239 (Some(a), Some(b)) if b > a => ((b - a) / 1000).max(1),
1240 _ => 1,
1241 };
1242 let run = self.run_row(&job.run_id).await?;
1243 let repo = run.as_ref().map(|r| r.repo.clone());
1244 let recorded: Result<Outcome<bool>> = g1t_kit::call(
1245 &self.billing,
1246 "record_sandbox",
1247 &RecordSandboxArgs {
1248 workspace: job.namespace.to_lowercase(),
1249 seconds: u32::try_from(seconds).unwrap_or(u32::MAX),
1250 description: format!(
1251 "{} in {} on the self-hosted runner {}",
1252 job.name,
1253 repo.clone().unwrap_or_default(),
1254 job.runner_name.clone().unwrap_or_default()
1255 ),
1256 repo,
1257 reference: format!("selfhosted/{}/{}", job.id, started),
1258 kind: Some(ComputeKind::Workflow),
1259 cpu_seconds: None,
1260 reservation_id: None,
1261 self_hosted: true,
1262 instance: None,
1263 },
1264 )
1265 .await;
1266 if let Err(error) = recorded {
1267 worker::console_error!("actions: self-hosted time not recorded for {}: {error}", job.id);
1268 }
1269 Ok(())
1270 }
1271
1272 /// `runner` and `RUNNER_*` for a job on a self-hosted runner.
1273 pub async fn runner_context_for(&self, runner_id: &str, variables: &mut Map<String, Value>) -> Result<Value> {
1274 let mut context = g1t_actions::events::runner_context();
1275 let Some(runner) = self.runner_by_id(runner_id).await? else { return Ok(context) };
1276 let os = match runner.os.as_str() {
1277 "macos" => "macOS",
1278 "windows" => "Windows",
1279 _ => "Linux",
1280 };
1281 let arch = if runner.arch == "arm64" { "ARM64" } else { "X64" };
1282 context["name"] = json!(runner.name);
1283 context["os"] = json!(os);
1284 context["arch"] = json!(arch);
1285 context["environment"] = json!("self-hosted");
1286 {
1287 let vars = variables;
1288 vars.insert("RUNNER_NAME".into(), json!(runner.name));
1289 vars.insert("RUNNER_OS".into(), json!(os));
1290 vars.insert("RUNNER_ARCH".into(), json!(arch));
1291 vars.insert("RUNNER_ENVIRONMENT".into(), json!("self-hosted"));
1292 }
1293 Ok(context)
1294 }
1295
1296 // --- Agent work, for the runner service ----------------------------------
1297
1298 /// `runner_route`: the labels g1t's own work in `repo` runs on, when
1299 /// its workspace (or the repository) sends it to self-hosted runners.
1300 pub async fn runner_route(&self, a: RouteArgs) -> Result<Option<Vec<String>>> {
1301 let repo_id = match self.repo_by_path(&a.repo).await? {
1302 Some(repo) => Some(repo.id),
1303 None => None,
1304 };
1305 let settings = self.effective_runner_settings(&a.workspace, repo_id.as_deref()).await?;
1306 Ok(settings.agents_on_self_hosted.then_some(settings.agent_labels))
1307 }
1308
1309 async fn repo_by_path(&self, path: &RepoPath) -> Result<Option<Repo>> {
1310 let Some(actor) = self.workspace_actor(&path.namespace).await? else { return Ok(None) };
1311 self.visible_repo(path, &Some(actor)).await
1312 }
1313
1314 /// `enqueue_task`.
1315 pub async fn enqueue_task(&self, a: EnqueueTaskArgs) -> Result<Outcome<String>> {
1316 let Some(sealer) = &self.sealer else {
1317 return Ok(fail(FailureCode::Conflict, "Agent work cannot be handed to self-hosted runners yet: g1t's key for it is not set."));
1318 };
1319 let id = new_id("rtk", now_ms());
1320 let sealed = sealer.seal(&serde_json::to_string(&a.env)?, &id);
1321 let repo_id = self.repo_by_path(&a.repo).await?.map(|r| r.id);
1322 let labels = match model::runner_labels(&a.labels, model::SELF_HOSTED, model::SELF_HOSTED) {
1323 Ok(labels) => labels,
1324 Err(problem) => return Ok(fail(FailureCode::Invalid, problem)),
1325 };
1326 let at = now();
1327 let inserted = self
1328 .db
1329 .prepare(
1330 "INSERT OR IGNORE INTO runner_tasks (id, sandbox, workspace, repo_id, repo, kind, title, labels, env, timeout_minutes, status, created_at)
1331 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?) RETURNING id",
1332 )
1333 .bind(&[
1334 id.as_str().into(),
1335 a.sandbox.as_str().into(),
1336 a.workspace.to_lowercase().into(),
1337 optional(repo_id.as_deref()),
1338 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
1339 a.kind.as_str().into(),
1340 a.title.chars().take(300).collect::<String>().into(),
1341 serde_json::to_string(&labels)?.into(),
1342 sealed.into(),
1343 a.timeout_minutes.max(1).into(),
1344 at.into(),
1345 ])?
1346 .first::<Value>(None)
1347 .await?;
1348 if inserted.is_none() {
1349 return Ok(fail(FailureCode::Conflict, "That sandbox already handed work to a self-hosted runner."));
1350 }
1351 Ok(Outcome::Ok(id))
1352 }
1353
1354 /// `cancel_task`: the sandbox stopped it. Its runner hears on its next
1355 /// poll; the sandbox does its own ending.
1356 pub async fn cancel_task(&self, a: CancelTaskArgs) -> Result<Outcome<bool>> {
1357 let done = self
1358 .db
1359 .prepare(
1360 "UPDATE runner_tasks SET status = 'completed', exit_code = -1, reason = ?, env = NULL, finished_at = ?
1361 WHERE sandbox = ? AND status != 'completed' RETURNING id",
1362 )
1363 .bind(&[optional(a.reason.as_deref()), now().into(), a.sandbox.as_str().into()])?
1364 .first::<Value>(None)
1365 .await?;
1366 Ok(Outcome::Ok(done.is_some()))
1367 }
1368
1369 /// Ends a task and tells the sandbox waiting on it, once.
1370 async fn end_task(&self, task: &TaskRow, exit_code: i32, reason: &str) -> Result<()> {
1371 let ended = self
1372 .db
1373 .prepare(
1374 "UPDATE runner_tasks SET status = 'completed', exit_code = ?, reason = ?, env = NULL, finished_at = ?
1375 WHERE id = ? AND status != 'completed' RETURNING id",
1376 )
1377 .bind(&[exit_code.into(), optional(Some(reason).filter(|r| !r.is_empty())), now().into(), task.id.as_str().into()])?
1378 .first::<Value>(None)
1379 .await?;
1380 if ended.is_none() {
1381 return Ok(());
1382 }
1383 if let Some(runner_id) = &task.runner_id {
1384 self.db
1385 .prepare("UPDATE runners SET work_id = NULL, work_kind = NULL WHERE id = ? AND work_id = ?")
1386 .bind(&[runner_id.as_str().into(), task.id.as_str().into()])?
1387 .run()
1388 .await?;
1389 }
1390 let told: Result<Value> = g1t_kit::call(
1391 &self.runner,
1392 "task_ended",
1393 &json!({
1394 "sandbox": task.sandbox,
1395 "exitCode": exit_code,
1396 "reason": if reason.is_empty() { Value::Null } else { json!(reason) },
1397 "runner": task.runner_name,
1398 }),
1399 )
1400 .await;
1401 if let Err(error) = told {
1402 worker::console_error!("actions: the sandbox waiting on task {} was not told: {error}", task.id);
1403 }
1404 Ok(())
1405 }
1406
1407 /// `stuck_jobs`: jobs of the viewer's workspaces that have waited ten
1408 /// minutes or more for a self-hosted runner while none that could take
1409 /// them is online, for Mission control's Needs you.
1410 pub async fn stuck_jobs(&self, a: StuckJobsArgs) -> Result<Vec<StuckJob>> {
1411 let Some(viewer) = a.viewer else { return Ok(Vec::new()) };
1412 let at = now_ms();
1413 let mut out = Vec::new();
1414 for membership in viewer.workspaces.iter().take(20) {
1415 let workspace = membership.slug.to_lowercase();
1416 #[derive(Deserialize)]
1417 struct Waiting {
1418 id: String,
1419 name: String,
1420 run_id: String,
1421 labels: String,
1422 queued_at: Option<String>,
1423 repo: String,
1424 }
1425 let waiting = self
1426 .db
1427 .prepare(
1428 "SELECT jobs.id, jobs.name, jobs.run_id, jobs.labels, jobs.queued_at, runs.repo FROM jobs JOIN runs ON runs.id = jobs.run_id
1429 WHERE jobs.status = 'queued' AND jobs.labels IS NOT NULL AND lower(jobs.namespace) = ? AND jobs.queued_at < ?
1430 ORDER BY jobs.rowid LIMIT 20",
1431 )
1432 .bind(&[workspace.as_str().into(), rfc3339(at.saturating_sub(10 * 60 * 1000)).into()])?
1433 .all()
1434 .await?
1435 .results::<Waiting>()?;
1436 if waiting.is_empty() {
1437 continue;
1438 }
1439 let groups = self.group_rows(&workspace).await?;
1440 let online: Vec<RunnerRow> = self.runner_rows(&workspace).await?.into_iter().filter(|r| r.online(at)).collect();
1441 for job in waiting {
1442 let wanted = Wanted::from_stored(&serde_json::from_str::<Vec<String>>(&job.labels).unwrap_or_default());
1443 let served = online.iter().any(|runner| {
1444 let group = runner.group_id.as_ref().and_then(|id| groups.iter().find(|g| &g.id == id));
1445 wanted.matches(&runner.labels(), group.map(|g| g.name.as_str()))
1446 });
1447 if !served {
1448 out.push(StuckJob {
1449 id: job.id,
1450 name: job.name,
1451 run_id: job.run_id,
1452 repo: job.repo,
1453 labels: wanted.describe(),
1454 queued_at: job.queued_at.unwrap_or_default(),
1455 });
1456 }
1457 }
1458 }
1459 Ok(out)
1460 }
1461
1462 // --- Every minute ----------------------------------------------------------
1463
1464 /// Jobs and tasks that waited too long, tasks whose runner went quiet,
1465 /// runners offline for weeks, and registration tokens long expired.
1466 pub async fn sweep_runners(&self, at: u64) -> Result<()> {
1467 let before = |ms: u64| rfc3339(at.saturating_sub(ms));
1468 let waited = model::MAX_WAIT_HOURS * 60 * 60 * 1000;
1469 let stale = self
1470 .db
1471 .prepare("SELECT * FROM jobs WHERE status = 'queued' AND labels IS NOT NULL AND queued_at < ? LIMIT 50")
1472 .bind(&[before(waited).into()])?
1473 .all()
1474 .await?
1475 .results::<JobRow>()?;
1476 for job in stale {
1477 let stored: Vec<String> = job.labels.as_deref().and_then(|l| serde_json::from_str(l).ok()).unwrap_or_default();
1478 let why = format!(
1479 "No self-hosted runner with labels {} took it within {} hours.",
1480 Wanted::from_stored(&stored).describe(),
1481 model::MAX_WAIT_HOURS
1482 );
1483 self.finish_job(&job.id, "failure", Some(&why), None).await?;
1484 }
1485 let tasks = self
1486 .db
1487 .prepare(
1488 "SELECT * FROM runner_tasks WHERE (status = 'queued' AND created_at < ?) OR (status = 'in_progress' AND seen_at < ?) LIMIT 50",
1489 )
1490 .bind(&[before(waited).into(), before(SILENT_MS).into()])?
1491 .all()
1492 .await?
1493 .results::<TaskRow>()?;
1494 for task in tasks {
1495 let why = if task.status == "queued" {
1496 format!("No self-hosted runner with labels {} took it within {} hours.", Wanted::of(&serde_json::from_str::<Vec<String>>(&task.labels).unwrap_or_default(), None).describe(), model::MAX_WAIT_HOURS)
1497 } else {
1498 format!("The self-hosted runner {} stopped answering.", task.runner_name.clone().unwrap_or_default())
1499 };
1500 self.end_task(&task, 1, &why).await?;
1501 }
1502 let forgotten = self
1503 .db
1504 .prepare("SELECT * FROM runners WHERE COALESCE(last_seen_at, created_at) < ? LIMIT 50")
1505 .bind(&[before(FORGET_OFFLINE_MS).into()])?
1506 .all()
1507 .await?
1508 .results::<RunnerRow>()?;
1509 for row in forgotten {
1510 self.forget_runner(&row, "The runner was offline for 14 days.").await?;
1511 }
1512 // An ephemeral runner that took its job and then went quiet.
1513 let spent = self
1514 .db
1515 .prepare("SELECT * FROM runners WHERE ephemeral = 1 AND spent = 1 AND work_id IS NULL AND COALESCE(last_seen_at, created_at) < ? LIMIT 50")
1516 .bind(&[before(60 * 60 * 1000).into()])?
1517 .all()
1518 .await?
1519 .results::<RunnerRow>()?;
1520 for row in spent {
1521 self.forget_runner(&row, "The ephemeral runner finished its job.").await?;
1522 }
1523 self.db
1524 .prepare("DELETE FROM runner_registrations WHERE expires_at < ?")
1525 .bind(&[before(24 * 60 * 60 * 1000).into()])?
1526 .run()
1527 .await?;
1528 Ok(())
1529 }
1530}
1531
1532#[derive(Deserialize)]
1533struct IdRow {
1534 id: String,
1535}
1536
1537fn group_view(group: &GroupRow, runners: &[RunnerRow]) -> RunnerGroup {
1538 RunnerGroup {
1539 id: group.id.clone(),
1540 name: group.name.clone(),
1541 default: group.is_default != 0,
1542 repositories: group.repositories(),
1543 runners: runners.iter().filter(|r| r.group_id.as_deref() == Some(group.id.as_str())).count() as u32,
1544 updated_at: group.updated_at.clone(),
1545 }
1546}
1547
1548#[cfg(test)]
1549mod tests {
1550 use super::*;
1551
1552 fn group(repositories: &[&str]) -> GroupRow {
1553 GroupRow {
1554 id: "rng_1".into(),
1555 name: "Default".into(),
1556 is_default: 1,
1557 repositories: serde_json::to_string(&repositories).unwrap(),
1558 updated_at: "2026-10-06T00:00:00Z".into(),
1559 }
1560 }
1561
1562 #[test]
1563 fn being_seen_is_written_once_a_minute() {
1564 let at = 2_000_000_000_000;
1565 assert!(seen_is_stale(None, at));
1566 assert!(!seen_is_stale(Some(&rfc3339(at - 59_000)), at));
1567 assert!(seen_is_stale(Some(&rfc3339(at - 61_000)), at));
1568 // A runner written down at the last moment is still online until the
1569 // next write, a poll later.
1570 assert!(SEEN_RESOLUTION_MS + model::MAX_POLL_WAIT_MS < model::ONLINE_WITHIN_MS);
1571 assert!(SEEN_RESOLUTION_MS < crate::SILENT_MS);
1572 }
1573
1574 #[test]
1575 fn a_group_lets_its_repositories_use_it() {
1576 assert!(group(&[]).allows("acme/web"));
1577 assert!(group(&["web", "api"]).allows("acme/WEB"));
1578 assert!(group(&["web"]).allows("web"));
1579 assert!(!group(&["web"]).allows("acme/docs"));
1580 }
1581
1582 #[test]
1583 fn names_are_checked() {
1584 assert_eq!(valid_name(" build-01 ").unwrap(), "build-01");
1585 assert!(valid_name("").is_err());
1586 assert!(valid_name("has space").is_err());
1587 assert!(valid_name(&"x".repeat(65)).is_err());
1588 assert!(valid_group_name("GPU machines").is_ok());
1589 assert!(valid_group_name("").is_err());
1590 }
1591
1592 #[test]
1593 fn online_means_seen_in_the_last_ninety_seconds() {
1594 let row = RunnerRow {
1595 id: "rnr_1".into(),
1596 workspace: "acme".into(),
1597 repo_id: None,
1598 repo: None,
1599 group_id: None,
1600 name: "a".into(),
1601 labels: "[]".into(),
1602 os: "linux".into(),
1603 arch: "x64".into(),
1604 version: String::new(),
1605 ephemeral: 0,
1606 credential_hash: String::new(),
1607 previous_hash: None,
1608 rotated_at: String::new(),
1609 work_id: None,
1610 work_kind: None,
1611 spent: 0,
1612 last_seen_at: Some(rfc3339(1_000_000_000)),
1613 created_at: String::new(),
1614 created_by: None,
1615 };
1616 assert!(row.online(1_000_000_000 + 60_000));
1617 assert!(!row.online(1_000_000_000 + 120_000));
1618 assert_eq!(row.workspace, "acme");
1619 }
1620}