Skip to content

g1t/services/actions/src/runners.rs

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