Skip to content
1,669 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

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

This file's history is long; its oldest lines are credited to the oldest commit read.