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/apps/api/src/runners.rs

84 lines3,863 bytesCodeBlame
1//! What a self-hosted runner calls: registering with a registration token,
2//! then polling for work, saying how work ended and removing itself, with
3//! its own credential. Neither is a g1t access token, so these are served
4//! before anything reads one, and no other endpoint accepts either.
5//!
6//! | Route | Body |
7//! | --- | --- |
8//! | `POST /runners/register` | `token`, `name`, `labels`, `os`, `arch`, `version`, `ephemeral`, `group`, `replace` |
9//! | `POST /runners/{id}/poll` | `version`, `running`, `wait_ms` |
10//! | `POST /runners/{id}/finished` | `id`, `exit_code`, `reason` |
11//! | `POST /runners/{id}/remove` | nothing |
12//!
13//! The last three send `Authorization: Bearer g1tr_…`. See
14//! `g1t_contracts::runners` for what each answers.
15
16use g1t_contracts::Outcome;
17use g1t_kit::wire;
18use serde_json::{Value, json};
19use worker::{Request, Response, Result};
20
21use crate::operations::Services;
22use crate::{fail, failure, json_body};
23
24/// What is passed through as given: an agent task's environment.
25const AS_GIVEN: &[&str] = &["env"];
26
27fn credential(request: &Request) -> Result<String> {
28 let header = request.headers().get("authorization")?.unwrap_or_default();
29 Ok(match header.split_once(' ') {
30 Some((scheme, token)) if scheme.eq_ignore_ascii_case("bearer") => token.trim().to_owned(),
31 _ => String::new(),
32 })
33}
34
35async fn answer(services: &Services, method: &str, args: &Value) -> Result<Response> {
36 let answered: Outcome<Value> = g1t_kit::call(&services.actions, method, args).await?;
37 match answered {
38 Outcome::Ok(value) => Response::from_json(&wire::snake_case_keeping(value, AS_GIVEN)),
39 Outcome::Fail(refused) => failure(&refused),
40 }
41}
42
43/// Serves `path` if it is one of the runner's routes.
44pub async fn handle(request: &mut Request, services: &Services, path: &str) -> Result<Option<Response>> {
45 let Some(rest) = path.strip_prefix("/runners/") else { return Ok(None) };
46 let body = json_body(request).await;
47 let body = if body.is_object() { body } else { json!({}) };
48 if rest == "register" {
49 let args = json!({
50 "token": body["token"].as_str().unwrap_or_default(),
51 "name": body["name"].as_str().unwrap_or_default(),
52 "labels": body["labels"].as_array().cloned().unwrap_or_default(),
53 "os": body["os"].as_str().unwrap_or_default(),
54 "arch": body["arch"].as_str().unwrap_or_default(),
55 "version": body["version"].as_str().unwrap_or_default(),
56 "ephemeral": body["ephemeral"].as_bool().unwrap_or(false),
57 "group": body["group"].as_str(),
58 "replace": body["replace"].as_bool().unwrap_or(false),
59 });
60 return Ok(Some(answer(services, "runner_register", &args).await?));
61 }
62 let Some((runner, action)) = rest.split_once('/') else {
63 return Ok(Some(fail(g1t_contracts::FailureCode::NotFound, "No such endpoint.")?));
64 };
65 let auth = json!({ "runner": runner, "credential": credential(request)? });
66 let mut args = auth;
67 let method = match action {
68 "poll" => {
69 args["version"] = json!(body["version"].as_str().unwrap_or_default());
70 args["running"] = json!(body["running"].as_array().cloned().unwrap_or_default());
71 args["wait_ms"] = json!(body["wait_ms"].as_u64().unwrap_or(0));
72 "runner_poll"
73 }
74 "finished" => {
75 args["id"] = json!(body["id"].as_str().unwrap_or_default());
76 args["exit_code"] = json!(body["exit_code"].as_i64().unwrap_or(1));
77 args["reason"] = body["reason"].clone();
78 "runner_finished"
79 }
80 "remove" => "runner_remove_self",
81 _ => return Ok(Some(fail(g1t_contracts::FailureCode::NotFound, "No such endpoint.")?)),
82 };
83 Ok(Some(answer(services, method, &args).await?))
84}