flagon-io/g1t

public

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

g1t/services/actions/src/lib.rs

301 lines13,059 bytesCodeBlame
1//! The actions service: a repository's GitHub Actions workflows, run on
2//! g1t as they are. See `g1t_contracts::actions` for the methods and
3//! `g1t_actions` for how workflows, expressions and filters are read.
4//!
5//! - [`sync`] reads workflow files, at the default branch for the list and
6//! at an event's own commit for its runs.
7//! - [`trigger`] turns events, schedules and manual runs into runs.
8//! - [`plan`] moves a run's jobs along: each waits for the jobs it needs,
9//! is skipped or expanded into its matrix, queued, started in a sandbox
10//! when the workspace has room, and finished by the sandbox's report.
11//! - [`payload`] builds the webhook-shaped `github.event`.
12//! - [`settings`] keeps secrets and variables.
13//! - [`cache`] lists `actions/cache` entries, which the API keeps in R2.
14//! - [`runners`] keeps self-hosted runners, hands them jobs (and agent
15//! work from the runner service) when they ask, and hears back.
16//!
17//! The service acts as the repository's workspace: it reads what the
18//! workspace can read, and a job's `GITHUB_TOKEN` is a short-lived token of
19//! the workspace's.
20
21mod cache;
22mod payload;
23mod plan;
24mod rename;
25mod runners;
26mod settings;
27mod sync;
28mod trigger;
29mod views;
30
31use g1t_contracts::access::{self, Capability};
32use g1t_contracts::events::Event;
33use g1t_contracts::identity::{SlugArgs, Workspace};
34use g1t_contracts::repos::{GetArgs, GetByIdArgs, Repo, RepoPath};
35use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, User, Viewer};
36use g1t_kit::{args, reply, rpc_method};
37use g1t_secrets::Sealer;
38use serde::Deserialize;
39use serde_json::Value;
40use worker::wasm_bindgen::JsValue;
41use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
42
43/// The most workflow files read from a repository.
44pub const MAX_WORKFLOWS: usize = 50;
45/// Jobs one workspace may have running at once; the rest wait their turn.
46pub const RUNNING_PER_WORKSPACE: u32 = 4;
47/// The longest a job may run, whatever its `timeout-minutes`.
48pub const MAX_TIMEOUT_MINUTES: u32 = 60;
49/// The longest a job on a self-hosted runner may run: the machine is the
50/// workspace's own, and its time costs nothing.
51pub const SELF_HOSTED_MAX_TIMEOUT_MINUTES: u32 = 24 * 60;
52/// A running job that has said nothing for this long is taken as lost.
53pub const SILENT_MS: u64 = 10 * 60 * 1000;
54pub const SITE: &str = "https://g1t.sh";
55pub const API: &str = "https://api.g1t.sh";
56
57#[derive(Deserialize)]
58struct Count {
59 n: u32,
60}
61
62pub fn optional(value: Option<&str>) -> JsValue {
63 value.map_or(JsValue::NULL, JsValue::from)
64}
65
66pub fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
67 Outcome::fail(code, message)
68}
69
70/// `owner/name` as a path.
71pub fn repo_path(full_name: &str) -> RepoPath {
72 let (namespace, name) = full_name.split_once('/').unwrap_or((full_name, ""));
73 RepoPath {
74 namespace: namespace.to_owned(),
75 name: name.to_owned(),
76 }
77}
78
79pub struct Actions {
80 db: D1Database,
81 repos: Fetcher,
82 work: Fetcher,
83 identity: Fetcher,
84 runner: Fetcher,
85 events: Fetcher,
86 /// Projects: a repository's secrets and variables belong to its project.
87 projects: Fetcher,
88 /// Billing: self-hosted runners' time, recorded at $0, and the
89 /// cache's storage.
90 billing: Fetcher,
91 /// Where the API keeps cache entries, for deleting evicted ones.
92 cache: Option<worker::Bucket>,
93 /// Seals secrets; absent until `ACTIONS_KEY` is set, when secrets
94 /// cannot be saved.
95 sealer: Option<Sealer>,
96}
97
98impl Actions {
99 fn new(env: &Env) -> Result<Self> {
100 Ok(Actions {
101 db: env.d1("DB")?,
102 repos: env.service("REPOS")?,
103 work: env.service("WORK")?,
104 identity: env.service("IDENTITY")?,
105 runner: env.service("RUNNER")?,
106 events: env.service("EVENTS")?,
107 projects: env.service("PROJECTS")?,
108 billing: env.service("BILLING")?,
109 cache: env.bucket("ACTIONS_CACHE").ok(),
110 sealer: env.secret("ACTIONS_KEY").ok().and_then(|key| Sealer::new(&key.to_string())),
111 })
112 }
113
114 /// The workspace itself, as the service acts.
115 async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> {
116 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?;
117 Ok(workspace.map(|workspace| User {
118 id: workspace.id,
119 username: workspace.slug.clone(),
120 kind: PrincipalKind::Workspace,
121 verified: true,
122 workspaces: vec![Membership::member(workspace.slug)],
123 ..User::default()
124 }))
125 }
126
127 /// The repository, if the viewer may see it and it is not a pull
128 /// request's working copy.
129 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
130 let found: Outcome<Repo> = g1t_kit::call(
131 &self.repos,
132 "get",
133 &GetArgs {
134 path: path.clone(),
135 viewer: viewer.clone(),
136 },
137 )
138 .await?;
139 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
140 }
141
142 /// A repository by id, as its workspace sees it.
143 async fn repo_by_id(&self, id: &str) -> Result<Option<(Repo, User)>> {
144 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &g1t_contracts::repos::PathByIdArgs { id: id.to_owned() }).await?;
145 let Some(path) = path else { return Ok(None) };
146 let Some(actor) = self.workspace_actor(&path.namespace).await? else {
147 return Ok(None);
148 };
149 let found: Outcome<Repo> = g1t_kit::call(
150 &self.repos,
151 "get_by_id",
152 &GetByIdArgs {
153 id: id.to_owned(),
154 viewer: Some(actor.clone()),
155 },
156 )
157 .await?;
158 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()).map(|repo| (repo, actor)))
159 }
160
161 /// The repository at `path`, when `actor` may do `capability` in it:
162 /// not found when they cannot read it, forbidden when their role falls
163 /// short. Agents never may: people and tokens run and change workflows.
164 async fn may(&self, actor: &User, path: &RepoPath, capability: Capability) -> Result<Outcome<Repo>> {
165 if actor.kind == PrincipalKind::Agent {
166 return Ok(fail(FailureCode::Forbidden, "An agent cannot do that. Ask a person."));
167 }
168 let Some(repo) = self.visible_repo(path, &Some(actor.clone())).await? else {
169 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
170 };
171 if !access::can(Some(actor), &repo, capability) {
172 return Ok(fail(
173 FailureCode::Forbidden,
174 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
175 ));
176 }
177 Ok(Outcome::Ok(repo))
178 }
179}
180
181/// Unwraps an `Outcome`, or returns its failure from the enclosing method.
182#[macro_export]
183macro_rules! check {
184 ($outcome:expr) => {
185 match $outcome {
186 g1t_contracts::Outcome::Ok(value) => value,
187 g1t_contracts::Outcome::Fail(refused) => return Ok(g1t_contracts::Outcome::Fail(refused)),
188 }
189 };
190}
191
192#[event(fetch)]
193async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
194 let Some(method) = rpc_method(&request) else {
195 return Response::error("Not found", 404);
196 };
197 let body: Value = request.json().await?;
198 let service = Actions::new(&env)?;
199 match method.as_str() {
200 "workflows" => reply(&service.workflows(args(body)?).await?),
201 "runs" => reply(&service.runs(args(body)?).await?),
202 "run" => reply(&service.run(args(body)?).await?),
203 "logs" => reply(&service.logs(args(body)?).await?),
204 "dispatch" => reply(&service.dispatch(args(body)?).await?),
205 "merge_group" => reply(&service.merge_group(args(body)?).await?),
206 "cancel" => reply(&service.cancel(args(body)?).await?),
207 "rerun" => reply(&service.rerun(args(body)?).await?),
208 "set_workflow_enabled" => reply(&service.set_workflow_enabled(args(body)?).await?),
209 "settings" => reply(&service.settings(args(body)?).await?),
210 "set_setting" => reply(&service.set_setting(args(body)?).await?),
211 "delete_setting" => reply(&service.delete_setting(args(body)?).await?),
212 "resolve_settings" => reply(&service.resolve_settings(args(body)?).await?),
213 "job_spec" => reply(&service.job_spec(args(body)?).await?),
214 "job_auth" => reply(&service.job_auth(args(body)?).await?),
215 "job_report" => reply(&service.job_report(args(body)?).await?),
216 // actions/cache, through the API with the job's token.
217 "cache_lookup" => reply(&service.cache_lookup(args(body)?).await?),
218 "cache_reserve" => reply(&service.cache_reserve(args(body)?).await?),
219 "cache_commit" => reply(&service.cache_commit(args(body)?).await?),
220 "cache_abort" => reply(&service.cache_abort(args(body)?).await?),
221 // Self-hosted runners: people's side.
222 "runners" => reply(&service.runners(args(body)?).await?),
223 "create_registration_token" => reply(&service.create_registration_token(args(body)?).await?),
224 "remove_runner" => reply(&service.remove_runner(args(body)?).await?),
225 "runner_groups" => reply(&service.runner_groups(args(body)?).await?),
226 "set_runner_group" => reply(&service.set_runner_group(args(body)?).await?),
227 "delete_runner_group" => reply(&service.delete_runner_group(args(body)?).await?),
228 "runner_settings" => reply(&service.runner_settings(args(body)?).await?),
229 "set_runner_settings" => reply(&service.set_runner_settings(args(body)?).await?),
230 // The runner's own side, through the API with its credential.
231 "runner_register" => reply(&service.runner_register(args(body)?).await?),
232 "runner_poll" => reply(&service.runner_poll(args(body)?).await?),
233 "runner_finished" => reply(&service.runner_finished(args(body)?).await?),
234 "runner_remove_self" => reply(&service.runner_remove_self(args(body)?).await?),
235 // Agent work, from the runner service.
236 "runner_route" => reply(&service.runner_route(args(body)?).await?),
237 "stuck_jobs" => reply(&service.stuck_jobs(args(body)?).await?),
238 "enqueue_task" => reply(&service.enqueue_task(args(body)?).await?),
239 "cancel_task" => reply(&service.cancel_task(args(body)?).await?),
240 _ => Response::error("Unknown method", 404),
241 }
242}
243
244/// Events from the bus, on this service's own queue.
245#[event(queue)]
246async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
247 let service = Actions::new(&env)?;
248 for message in batch.messages()? {
249 // A workspace renamed: its rows move to the slug it has now.
250 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), rename::STATEMENTS).await? {
251 message.ack();
252 continue;
253 }
254 // A repository renamed or transferred: its rows follow its new path.
255 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), rename::TRANSFERRED).await? {
256 message.ack();
257 continue;
258 }
259 // A workspace deleted: what it kept for itself goes.
260 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), rename::DELETED).await? {
261 message.ack();
262 continue;
263 }
264 // A repository purged: every row kept for it goes.
265 if g1t_kit::lifecycle::on_purged(&env.d1("DB")?, message.body(), rename::PURGED).await? {
266 message.ack();
267 continue;
268 }
269 // A repository deleted or archived: its runs stop.
270 if let Some(repo_id) = plan::stops_runs(message.body()) {
271 if let Err(error) = service.stop_runs(&repo_id).await {
272 worker::console_error!("actions: event {} failed: {error}", message.body().id);
273 message.retry();
274 continue;
275 }
276 message.ack();
277 continue;
278 }
279 if let Err(error) = service.on_event(message.body()).await {
280 worker::console_error!("actions: event {} failed: {error}", message.body().id);
281 message.retry();
282 continue;
283 }
284 message.ack();
285 }
286 Ok(())
287}
288
289/// Every minute: schedules that fire, jobs waiting for room, and jobs
290/// whose sandbox went quiet.
291#[event(scheduled)]
292async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
293 match Actions::new(&env) {
294 Ok(service) => {
295 if let Err(error) = service.on_minute(g1t_kit::now_ms()).await {
296 worker::console_error!("actions: the sweep failed: {error}");
297 }
298 }
299 Err(error) => worker::console_error!("actions: could not start: {error}"),
300 }
301}