g1t/crates/runner/src/selfhosted/mod.rs

526 lines21,011 bytesCodeBlame
1//! `g1t-runner` on your own machine: a self-hosted runner for g1t.
2//!
3//! ```text
4//! g1t-runner register --url https://g1t.sh --token g1trt_… [--name build-01]
5//! [--labels gpu,cuda] [--group Default] [--ephemeral] [--work-dir DIR]
6//! [--no-docker] [--image IMAGE] [--agent-image IMAGE] [--replace]
7//! g1t-runner run [--once]
8//! g1t-runner service install|start|stop|status|uninstall
9//! g1t-runner remove
10//! g1t-runner update
11//! g1t-runner version
12//! ```
13//!
14//! It only ever calls out, to `api.g1t.sh` over HTTPS: it polls for work
15//! (a poll waits up to 20 seconds for some, and is the runner's heartbeat),
16//! runs what it is given with the same harness g1t's sandboxes run (this
17//! program in another mode, see `exec`), and says how it ended. Nothing
18//! listens on the machine. Every command takes `--dir` (default
19//! `~/.g1t-runner`, or `G1T_RUNNER_DIR`), where its configuration is kept.
20
21mod api;
22mod config;
23mod exec;
24mod service;
25mod update;
26
27use std::path::{Path, PathBuf};
28use std::sync::atomic::{AtomicBool, Ordering};
29use std::time::{Duration, Instant};
30
31use anyhow::{Context, Result, bail};
32use serde_json::json;
33
34use api::{Api, Failure};
35use config::Config;
36
37pub const VERSION: &str = env!("CARGO_PKG_VERSION");
38
39/// Set to stop the loop between polls: a Windows service's stop request.
40pub static STOP: AtomicBool = AtomicBool::new(false);
41
42/// How long a busy runner waits between heartbeats.
43const HEARTBEAT: Duration = Duration::from_secs(10);
44/// How long an idle poll waits for work.
45const IDLE_WAIT_MS: u64 = 20_000;
46/// How often an idle runner checks for a new version.
47const UPDATE_EVERY: Duration = Duration::from_secs(6 * 60 * 60);
48
49const COMMANDS: &[&str] = &["register", "register-and-run", "run", "service", "remove", "update", "version", "--version", "help", "--help", "-h"];
50
51/// Whether `args` (without the program) are a command of this one, rather
52/// than the harness's `MODE`.
53pub fn is_command(args: &[String]) -> bool {
54 args.first().is_some_and(|first| COMMANDS.contains(&first.as_str()))
55}
56
57pub fn log(message: &str) {
58 let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map(|d| d.as_millis() as u64).unwrap_or(0);
59 println!("{} {message}", g1t_now(now));
60}
61
62/// RFC 3339 for the log, without pulling in a date library.
63fn g1t_now(ms: u64) -> String {
64 let seconds = ms / 1000;
65 let (days, rest) = (seconds / 86_400, seconds % 86_400);
66 let z = days as i64 + 719_468;
67 let era = z.div_euclid(146_097);
68 let doe = z.rem_euclid(146_097);
69 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
70 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
71 let mp = (5 * doy + 2) / 153;
72 let day = doy - (153 * mp + 2) / 5 + 1;
73 let month = if mp < 10 { mp + 3 } else { mp - 9 };
74 let year = yoe + era * 400 + i64::from(month <= 2);
75 format!("{year:04}-{month:02}-{day:02}T{:02}:{:02}:{:02}Z", rest / 3600, rest % 3600 / 60, rest % 60)
76}
77
78/// `--name value` and `--flag` options.
79struct Options {
80 values: Vec<(String, String)>,
81 flags: Vec<String>,
82 rest: Vec<String>,
83}
84
85impl Options {
86 fn parse(args: &[String], takes_value: &[&str]) -> Result<Options> {
87 let mut options = Options { values: Vec::new(), flags: Vec::new(), rest: Vec::new() };
88 let mut i = 0;
89 while i < args.len() {
90 let arg = &args[i];
91 if let Some(name) = arg.strip_prefix("--") {
92 if let Some((name, value)) = name.split_once('=') {
93 options.values.push((name.to_owned(), value.to_owned()));
94 } else if takes_value.contains(&name) {
95 let Some(value) = args.get(i + 1) else { bail!("--{name} needs a value") };
96 options.values.push((name.to_owned(), value.clone()));
97 i += 1;
98 } else {
99 options.flags.push(name.to_owned());
100 }
101 } else {
102 options.rest.push(arg.clone());
103 }
104 i += 1;
105 }
106 Ok(options)
107 }
108
109 fn value(&self, name: &str) -> Option<&str> {
110 self.values.iter().rev().find(|(n, _)| n == name).map(|(_, v)| v.as_str())
111 }
112
113 fn flag(&self, name: &str) -> bool {
114 self.flags.iter().any(|f| f == name)
115 }
116
117 /// For `register-and-run` in a container: each option may come from
118 /// `G1T_RUNNER_<NAME>` instead (`G1T_RUNNER_TOKEN`, `G1T_RUNNER_LABELS`,
119 /// `G1T_RUNNER_EPHEMERAL=1`), so a Kubernetes secret can hold the token.
120 fn with_environment(mut self) -> Options {
121 let env_name = |name: &str| format!("G1T_RUNNER_{}", name.to_ascii_uppercase().replace('-', "_"));
122 for name in ["url", "api", "token", "name", "labels", "group", "work-dir", "image", "agent-image", "harness"] {
123 if self.value(name).is_none()
124 && let Ok(value) = std::env::var(env_name(name))
125 && !value.is_empty()
126 {
127 self.values.push((name.to_owned(), value));
128 }
129 }
130 for name in ["ephemeral", "no-docker", "replace", "no-auto-update"] {
131 if !self.flag(name) && std::env::var(env_name(name)).is_ok_and(|v| matches!(v.as_str(), "1" | "true" | "yes")) {
132 self.flags.push(name.to_owned());
133 }
134 }
135 self
136 }
137}
138
139const HELP: &str = "g1t-runner: a self-hosted runner for g1t.
140
141 g1t-runner register --url https://g1t.sh --token <registration token>
142 [--name NAME] [--labels a,b] [--group GROUP] [--ephemeral]
143 [--work-dir DIR] [--no-docker] [--image IMAGE] [--agent-image IMAGE]
144 [--api URL] [--harness PATH] [--replace] [--no-auto-update]
145 g1t-runner run [--once] poll for work and run it
146 g1t-runner register-and-run ... both, for containers; options may be
147 G1T_RUNNER_URL, _TOKEN, _LABELS, _EPHEMERAL…
148 g1t-runner service install|start|stop|status|uninstall
149 g1t-runner remove unregister this runner
150 g1t-runner update update to the newest release
151 g1t-runner version
152
153Every command takes --dir DIR (default ~/.g1t-runner, or G1T_RUNNER_DIR).
154Make a registration token under Settings, Runners on g1t.sh.
155Docs: https://docs.g1t.sh/guides/self-hosted-runners/";
156
157/// Runs a command. Returns the process's exit code.
158pub fn main(args: Vec<String>) -> i32 {
159 let takes = ["url", "api", "token", "name", "labels", "group", "work-dir", "image", "agent-image", "harness", "dir"];
160 let options = match Options::parse(&args[1..], &takes) {
161 Ok(options) => options,
162 Err(error) => {
163 eprintln!("g1t-runner: {error:#}");
164 return 2;
165 }
166 };
167 let folder = config::folder(options.value("dir"));
168 let outcome = match args[0].as_str() {
169 "register" => register(&options, &folder),
170 // A container's one command: register unless it already is, then run.
171 "register-and-run" => {
172 let options = options.with_environment();
173 let ready = if config::load(&folder).is_ok() { Ok(()) } else { register(&options, &folder) };
174 ready.and_then(|()| run(&folder, options.flag("once")))
175 }
176 "run" => run(&folder, options.flag("once")),
177 "service" => match options.rest.first().map(String::as_str) {
178 Some(action) => config::load(&folder).and_then(|config| service::main(action, &config, &folder)),
179 None => Err(anyhow::anyhow!("service takes install, uninstall, start, stop or status")),
180 },
181 "remove" => remove(&folder),
182 "update" => config::load(&folder).and_then(|config| match update::update(&config)? {
183 Some(version) => {
184 log(&format!("Updated to {version}. Restart the runner (or its service) to use it."));
185 Ok(())
186 }
187 None => {
188 log(&format!("{VERSION} is the newest release."));
189 Ok(())
190 }
191 }),
192 "version" | "--version" => {
193 println!("g1t-runner {VERSION} ({})", update::platform());
194 Ok(())
195 }
196 _ => {
197 println!("{HELP}");
198 Ok(())
199 }
200 };
201 match outcome {
202 Ok(()) => 0,
203 Err(error) => {
204 eprintln!("g1t-runner: {error:#}");
205 1
206 }
207 }
208}
209
210fn hostname() -> String {
211 for name in ["COMPUTERNAME", "HOSTNAME"] {
212 if let Ok(value) = std::env::var(name).map(|v| v.trim().to_owned())
213 && !value.is_empty()
214 {
215 return value;
216 }
217 }
218 std::process::Command::new("hostname")
219 .output()
220 .ok()
221 .map(|out| String::from_utf8_lossy(&out.stdout).trim().to_owned())
222 .filter(|name| !name.is_empty())
223 .unwrap_or_else(|| "runner".to_owned())
224}
225
226/// A runner's name from a machine's: what g1t accepts.
227fn runner_name(given: &str) -> String {
228 let name: String = given
229 .chars()
230 .map(|c| if c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.') { c } else { '-' })
231 .take(64)
232 .collect();
233 if name.is_empty() { "runner".into() } else { name }
234}
235
236fn register(options: &Options, folder: &Path) -> Result<()> {
237 let Some(url) = options.value("url") else { bail!("--url is needed: the g1t you register with, such as https://g1t.sh") };
238 let Some(token) = options.value("token") else { bail!("--token is needed: make a registration token under Settings, Runners") };
239 let url = url.trim_end_matches('/').to_owned();
240 let api = options.value("api").map(|a| a.trim_end_matches('/').to_owned()).unwrap_or_else(|| config::api_for(&url));
241 let name = runner_name(options.value("name").map(str::to_owned).unwrap_or_else(hostname).as_str());
242 let labels: Vec<String> = options
243 .value("labels")
244 .unwrap_or_default()
245 .split(',')
246 .map(|l| l.trim().to_owned())
247 .filter(|l| !l.is_empty())
248 .collect();
249 let docker = !options.flag("no-docker");
250 if docker && !exec::docker_ready() {
251 log("Docker does not answer here. Jobs will fail until it does; register with --no-docker to run them on this machine instead.");
252 }
253 let existing = config::load(folder).ok();
254 if existing.is_some() && !options.flag("replace") {
255 bail!("{} already holds a registered runner: run `g1t-runner remove` first, or register with --replace or another --dir", folder.display());
256 }
257 // The OS jobs run on: in Docker that is Linux, whatever the machine
258 // (Docker Desktop runs Linux containers on macOS and Windows).
259 let os = match std::env::consts::OS {
260 _ if docker => "linux",
261 "macos" => "macos",
262 "windows" => "windows",
263 _ => "linux",
264 };
265 let arch = if std::env::consts::ARCH == "aarch64" { "arm64" } else { "x64" };
266 let registered = Api::register(
267 &api,
268 &json!({
269 "token": token,
270 "name": name,
271 "labels": labels,
272 "os": os,
273 "arch": arch,
274 "version": VERSION,
275 "ephemeral": options.flag("ephemeral"),
276 "group": options.value("group"),
277 "replace": options.flag("replace"),
278 }),
279 )?;
280 let work_dir = options.value("work-dir").map(PathBuf::from).unwrap_or_else(|| folder.join("work"));
281 let config = Config {
282 url,
283 api,
284 runner: registered.runner.id.clone(),
285 name: registered.runner.name.clone(),
286 credential: registered.credential,
287 workspace: registered.runner.workspace.clone(),
288 repo: registered.runner.repo.clone(),
289 group: registered.runner.group.clone(),
290 labels: registered.runner.labels.clone(),
291 ephemeral: options.flag("ephemeral"),
292 work_dir,
293 docker,
294 image: options.value("image").map(str::to_owned),
295 agent_image: options.value("agent-image").map(str::to_owned),
296 harness: options.value("harness").map(PathBuf::from),
297 auto_update: !options.flag("no-auto-update"),
298 };
299 config::save(folder, &config)?;
300 let place = match &config.repo {
301 Some(repo) => format!("for {repo}"),
302 None => format!("in {}{}", config.workspace, config.group.as_ref().map(|g| format!(", group {g}")).unwrap_or_default()),
303 };
304 log(&format!("Registered {} {place}, with labels {}.", config.name, config.labels.join(", ")));
305 log("Start it with `g1t-runner run`, or `g1t-runner service install` to keep it running.");
306 Ok(())
307}
308
309fn remove(folder: &Path) -> Result<()> {
310 let config = config::load(folder)?;
311 let api = Api { base: config.api.clone(), runner: config.runner.clone(), credential: config.credential.clone() };
312 match api.remove() {
313 Ok(()) | Err(Failure::Refused(401, _)) | Err(Failure::Refused(404, _)) => {}
314 Err(failure) => bail!("{failure}; the runner is still registered"),
315 }
316 config::forget(folder);
317 log(&format!("Removed {}.", config.name));
318 Ok(())
319}
320
321/// The loop, for a Windows service: its exit code.
322#[cfg_attr(not(windows), allow(dead_code))]
323pub fn serve_from_service() -> i32 {
324 let folder = config::folder(None);
325 let folder = std::env::args()
326 .collect::<Vec<_>>()
327 .windows(2)
328 .find(|pair| pair[0] == "--dir")
329 .map(|pair| PathBuf::from(&pair[1]))
330 .unwrap_or(folder);
331 match run(&folder, false) {
332 Ok(()) => 0,
333 Err(error) => {
334 eprintln!("g1t-runner: {error:#}");
335 1
336 }
337 }
338}
339
340/// Sleeps up to `duration`, waking early when `stop` says so.
341fn nap(duration: Duration, mut stop: impl FnMut() -> bool) {
342 let until = Instant::now() + duration;
343 while Instant::now() < until {
344 if STOP.load(Ordering::SeqCst) || stop() {
345 return;
346 }
347 std::thread::sleep(Duration::from_millis(250));
348 }
349}
350
351fn run(folder: &Path, once: bool) -> Result<()> {
352 let mut config = config::load(folder)?;
353 let mut api = Api { base: config.api.clone(), runner: config.runner.clone(), credential: config.credential.clone() };
354 std::fs::create_dir_all(&config.work_dir).with_context(|| format!("could not make {}", config.work_dir.display()))?;
355 log(&format!(
356 "g1t-runner {VERSION}: {} ({}) is listening for work from {}, {}.",
357 config.name,
358 config.labels.join(", "),
359 config.url,
360 if config.docker { "in Docker" } else { "on this machine" }
361 ));
362 let mut current: Option<exec::Running> = None;
363 let mut ran_one = false;
364 let mut backoff = Duration::from_secs(2);
365 let mut last_update_check = Instant::now().checked_sub(UPDATE_EVERY).unwrap_or_else(Instant::now);
366 let mut said_no_key = false;
367 loop {
368 if STOP.load(Ordering::SeqCst) {
369 if let Some(mut running) = current.take() {
370 log(&format!("Stopping: {} is cut short.", running.name));
371 running.stop();
372 let _ = api.finished(&running.id, 1, Some("The self-hosted runner was stopped."));
373 running.clean();
374 }
375 return Ok(());
376 }
377 // What was running has ended: say how.
378 if let Some(running) = current.as_mut()
379 && let Some(code) = running.ended()
380 {
381 let mut running = current.take().expect("checked above");
382 log(&format!("{} {} ended ({code}).", running.kind, running.name));
383 if let Err(failure) = api.finished(&running.id, code, None) {
384 log(&format!("Could not say how it ended: {failure}"));
385 }
386 running.clean();
387 ran_one = true;
388 if config.ephemeral || once {
389 if config.ephemeral {
390 let _ = api.remove();
391 config::forget(folder);
392 log("The ephemeral runner ran its job and removed itself.");
393 }
394 return Ok(());
395 }
396 }
397 // Idle: a newer release, now and then.
398 if current.is_none() && config.auto_update && last_update_check.elapsed() >= UPDATE_EVERY {
399 last_update_check = Instant::now();
400 if update::can_update() {
401 match update::update(&config) {
402 Ok(Some(version)) => {
403 log(&format!("Updated to {version}; starting it."));
404 return restart();
405 }
406 Ok(None) => {}
407 Err(error) => log(&format!("Could not check for a new version: {error:#}")),
408 }
409 } else if !said_no_key {
410 said_no_key = true;
411 log("This build cannot check releases, so it will not update itself.");
412 }
413 }
414 let running: Vec<String> = current.iter().map(|r| r.id.clone()).collect();
415 let wait = if current.is_some() { 0 } else { IDLE_WAIT_MS };
416 let poll = match api.poll(&running, wait) {
417 Ok(poll) => {
418 backoff = Duration::from_secs(2);
419 poll
420 }
421 Err(Failure::Refused(401, message)) => {
422 if let Some(mut running) = current.take() {
423 running.stop();
424 running.clean();
425 }
426 bail!("g1t no longer knows this runner ({message}). Register it again.");
427 }
428 Err(failure) => {
429 log(&format!("{failure}; trying again in {}s", backoff.as_secs()));
430 nap(backoff, || false);
431 backoff = (backoff * 2).min(Duration::from_secs(60));
432 continue;
433 }
434 };
435 if let Some(fresh) = poll.credential {
436 config.credential = fresh.clone();
437 api.credential = fresh;
438 config::save(folder, &config)?;
439 }
440 if poll.removed {
441 config::forget(folder);
442 log("This runner was removed from g1t; stopping.");
443 return Ok(());
444 }
445 if let Some(running) = current.as_mut()
446 && poll.cancel.contains(&running.id)
447 {
448 log(&format!("{} was cancelled; stopping it.", running.name));
449 running.stop();
450 }
451 if let Some(running) = current.as_mut()
452 && running.overdue()
453 {
454 log(&format!("{} ran past its time limit; stopping it.", running.name));
455 running.stop();
456 }
457 if let Some(work) = poll.assignment {
458 log(&format!("Took {} {} ({}).", work.kind, work.name, work.repo));
459 match exec::start(&config, folder, &work) {
460 Ok(running) => current = Some(running),
461 Err(error) => {
462 log(&format!("Could not start it: {error:#}"));
463 let _ = api.finished(&work.id, 1, Some(&format!("The self-hosted runner {} could not start it: {error:#}", config.name)));
464 if config.ephemeral || once {
465 return Ok(());
466 }
467 }
468 }
469 continue;
470 }
471 if once && ran_one {
472 return Ok(());
473 }
474 // Busy: a heartbeat every few seconds until it ends.
475 if current.is_some() {
476 let mut check = || current.as_mut().is_some_and(|r| r.ended().is_some());
477 nap(HEARTBEAT, &mut check);
478 }
479 }
480}
481
482/// Runs the updated program in this one's place, with the same arguments.
483fn restart() -> Result<()> {
484 let me = std::env::current_exe()?;
485 let status = std::process::Command::new(me).args(std::env::args().skip(1)).status()?;
486 std::process::exit(status.code().unwrap_or(1));
487}
488
489#[cfg(test)]
490mod tests {
491 use super::*;
492
493 fn args(items: &[&str]) -> Vec<String> {
494 items.iter().map(|s| s.to_string()).collect()
495 }
496
497 #[test]
498 fn commands_are_told_apart_from_the_harness() {
499 assert!(is_command(&args(&["register", "--url", "x"])));
500 assert!(is_command(&args(&["run"])));
501 assert!(!is_command(&args(&[])));
502 assert!(!is_command(&args(&["something"])));
503 }
504
505 #[test]
506 fn options_read_both_ways() {
507 let options = Options::parse(&args(&["--url", "https://g1t.sh", "--labels=gpu,cuda", "--ephemeral", "install"]), &["url", "labels"]).unwrap();
508 assert_eq!(options.value("url"), Some("https://g1t.sh"));
509 assert_eq!(options.value("labels"), Some("gpu,cuda"));
510 assert!(options.flag("ephemeral"));
511 assert_eq!(options.rest, vec!["install".to_owned()]);
512 assert!(Options::parse(&args(&["--url"]), &["url"]).is_err());
513 }
514
515 #[test]
516 fn a_machines_name_becomes_a_runners() {
517 assert_eq!(runner_name("Build Box #1"), "Build-Box--1");
518 assert_eq!(runner_name(""), "runner");
519 assert_eq!(runner_name(&"x".repeat(80)).len(), 64);
520 }
521
522 #[test]
523 fn the_log_says_when() {
524 assert_eq!(g1t_now(1_790_918_179_123), "2026-10-02T05:16:19Z");
525 }
526}