flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/actions/src/runners.rs

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