Skip to content
362 linesCodeBlameRaw
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//! - [`protection`] keeps a repository's workflow settings and its
14//! environments' protection rules, holds jobs at those rules until they
15//! pass, and lets pull requests' runs from outside wait for approval.
16//! - [`cache`] lists `actions/cache` entries, which the API keeps in R2.
17//! - [`runners`] keeps self-hosted runners, hands them jobs (and agent
18//! work from the runner service) when they ask, and hears back.
19//!
20//! The service acts as the repository's workspace: it reads what the
21//! workspace can read, and a job's `GITHUB_TOKEN` is a token of the
22//! workspace's that reaches the job's repository only, with the scopes its
23//! `permissions:` give it, revoked when the job ends.
24
25mod run_artifacts;
26mod cache;
27mod mirrored;
28mod payload;
29mod plan;
30mod protection;
31mod reach;
32mod rename;
33pub mod runtime;
34mod runners;
35mod settings;
36mod sync;
37mod trigger;
38mod views;
39
40use g1t_contracts::access::{self, Capability};
41use g1t_contracts::events::Event;
42use g1t_contracts::identity::{SlugArgs, Workspace};
43use g1t_contracts::repos::{GetArgs, GetByIdArgs, Repo, RepoPath};
44use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, User, Viewer};
45use g1t_kit::{args, reply, rpc_method};
46use g1t_secrets::Sealer;
47use serde::Deserialize;
48use serde_json::Value;
49use worker::wasm_bindgen::JsValue;
50use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
51
52/// The most workflow files read from a repository.
53pub const MAX_WORKFLOWS: usize = 50;
54/// Jobs one workspace may have running at once; the rest wait their turn.
55pub const RUNNING_PER_WORKSPACE: u32 = 4;
56/// The longest a job may run, whatever its `timeout-minutes`.
57pub const MAX_TIMEOUT_MINUTES: u32 = 60;
58/// The longest a job on a self-hosted runner may run: the machine is the
59/// workspace's own, and its time costs nothing.
60pub const SELF_HOSTED_MAX_TIMEOUT_MINUTES: u32 = 24 * 60;
61/// A running job that has said nothing for this long is taken as lost.
62pub const SILENT_MS: u64 = 10 * 60 * 1000;
63/// Schedules pause in a repository that has had no push for this long; the
64/// next push resumes them.
65pub const SCHEDULE_IDLE_MS: u64 = 60 * 24 * 60 * 60 * 1000;
66/// How long a workflow's schedule waits after billing refused to start a
67/// job of one of its scheduled runs, before it tries again.
68pub const SCHEDULE_REFUSED_MS: u64 = 60 * 60 * 1000;
69pub const SITE: &str = "https://g1t.sh";
70pub const API: &str = "https://api.g1t.sh";
71
72#[derive(Deserialize)]
73struct Count {
74 n: u32,
75}
76
77pub fn optional(value: Option<&str>) -> JsValue {
78 value.map_or(JsValue::NULL, JsValue::from)
79}
80
81pub fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
82 Outcome::fail(code, message)
83}
84
85/// `owner/name` as a path.
86pub fn repo_path(full_name: &str) -> RepoPath {
87 let (namespace, name) = full_name.split_once('/').unwrap_or((full_name, ""));
88 RepoPath {
89 namespace: namespace.to_owned(),
90 name: name.to_owned(),
91 }
92}
93
94pub struct Actions {
95 db: D1Database,
96 repos: Fetcher,
97 work: Fetcher,
98 identity: Fetcher,
99 runner: Fetcher,
100 events: Fetcher,
101 /// Projects: a repository's secrets and variables belong to its project.
102 projects: Fetcher,
103 /// Billing: self-hosted runners' time, recorded at $0, and the
104 /// cache's storage.
105 billing: Fetcher,
106 /// Deployments: a job with an `environment:` deploys to it, and its
107 /// run's deployment is recorded there (plan.rs `report_deployment`).
108 deployments: Fetcher,
109 /// Where the API keeps cache entries, for deleting evicted ones.
110 cache: Option<worker::Bucket>,
111 /// Seals secrets; absent until `ACTIONS_KEY` is set, when secrets
112 /// cannot be saved.
113 sealer: Option<Sealer>,
114 /// Signs the toolkit's blob links (runtime.rs), from `ACTIONS_KEY`.
115 blob_key: Option<String>,
116}
117
118impl Actions {
119 fn new(env: &Env) -> Result<Self> {
120 Ok(Actions {
121 db: env.d1("DB")?,
122 repos: env.service("REPOS")?,
123 work: env.service("WORK")?,
124 identity: env.service("IDENTITY")?,
125 runner: env.service("RUNNER")?,
126 events: env.service("EVENTS")?,
127 projects: env.service("PROJECTS")?,
128 billing: env.service("BILLING")?,
129 deployments: env.service("DEPLOYMENTS")?,
130 cache: env.bucket("ACTIONS_CACHE").ok(),
131 sealer: env.secret("ACTIONS_KEY").ok().and_then(|key| Sealer::new(&key.to_string())),
132 blob_key: env.secret("ACTIONS_KEY").ok().map(|key| runtime::blob_key(&key.to_string())),
133 })
134 }
135
136 /// The workspace itself, as the service acts.
137 async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> {
138 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?;
139 Ok(workspace.map(|workspace| User {
140 id: workspace.id,
141 username: workspace.slug.clone(),
142 kind: PrincipalKind::Workspace,
143 verified: true,
144 workspaces: vec![Membership::member(workspace.slug)],
145 ..User::default()
146 }))
147 }
148
149 /// The repository, if the viewer may see it and it is not a pull
150 /// request's working copy.
151 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
152 let found: Outcome<Repo> = g1t_kit::call(
153 &self.repos,
154 "get",
155 &GetArgs {
156 path: path.clone(),
157 viewer: viewer.clone(),
158 },
159 )
160 .await?;
161 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
162 }
163
164 /// A repository by id, as its workspace sees it.
165 async fn repo_by_id(&self, id: &str) -> Result<Option<(Repo, User)>> {
166 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &g1t_contracts::repos::PathByIdArgs { id: id.to_owned() }).await?;
167 let Some(path) = path else { return Ok(None) };
168 let Some(actor) = self.workspace_actor(&path.namespace).await? else {
169 return Ok(None);
170 };
171 let found: Outcome<Repo> = g1t_kit::call(
172 &self.repos,
173 "get_by_id",
174 &GetByIdArgs {
175 id: id.to_owned(),
176 viewer: Some(actor.clone()),
177 },
178 )
179 .await?;
180 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()).map(|repo| (repo, actor)))
181 }
182
183 /// The repository at `path`, when `actor` may do `capability` in it:
184 /// not found when they cannot read it, forbidden when their role falls
185 /// short. Agents never may: people and tokens run and change workflows.
186 async fn may(&self, actor: &User, path: &RepoPath, capability: Capability) -> Result<Outcome<Repo>> {
187 if actor.kind == PrincipalKind::Agent {
188 return Ok(fail(FailureCode::Forbidden, "An agent cannot do that. Ask a person."));
189 }
190 let Some(repo) = self.visible_repo(path, &Some(actor.clone())).await? else {
191 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
192 };
193 if !access::can(Some(actor), &repo, capability) {
194 return Ok(fail(
195 FailureCode::Forbidden,
196 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
197 ));
198 }
199 Ok(Outcome::Ok(repo))
200 }
201}
202
203/// Unwraps an `Outcome`, or returns its failure from the enclosing method.
204#[macro_export]
205macro_rules! check {
206 ($outcome:expr) => {
207 match $outcome {
208 g1t_contracts::Outcome::Ok(value) => value,
209 g1t_contracts::Outcome::Fail(refused) => return Ok(g1t_contracts::Outcome::Fail(refused)),
210 }
211 };
212}
213
214#[event(fetch)]
215async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
216 let Some(method) = rpc_method(&request) else {
217 return Response::error("Not found", 404);
218 };
219 let body: Value = request.json().await?;
220 let service = Actions::new(&env)?;
221 match method.as_str() {
222 "workflows" => reply(&service.workflows(args(body)?).await?),
223 "runs" => reply(&service.runs(args(body)?).await?),
224 "run" => reply(&service.run(args(body)?).await?),
225 "logs" => reply(&service.logs(args(body)?).await?),
226 "job_log_text" => reply(&service.job_log_text(args(body)?).await?),
227 "run_logs" => reply(&service.run_logs(args(body)?).await?),
228 "summaries" => reply(&service.summaries(args(body)?).await?),
229 "dispatch" => reply(&service.dispatch(args(body)?).await?),
230 "repository_dispatch" => reply(&service.repository_dispatch(args(body)?).await?),
231 "approve_run" => reply(&service.approve_run(args(body)?).await?),
232 "pending_deployments" => reply(&service.pending_deployments(args(body)?).await?),
233 "review_deployments" => reply(&service.review_deployments(args(body)?).await?),
234 "actions_settings" => reply(&service.actions_settings(args(body)?).await?),
235 "set_actions_settings" => reply(&service.set_actions_settings(args(body)?).await?),
236 "workspace_actions_settings" => reply(&service.workspace_actions_settings(args(body)?).await?),
237 "set_workspace_actions_settings" => reply(&service.set_workspace_actions_settings(args(body)?).await?),
238 "environments" => reply(&service.environments(args(body)?).await?),
239 "set_environment" => reply(&service.set_environment(args(body)?).await?),
240 "delete_environment" => reply(&service.delete_environment(args(body)?).await?),
241 "merge_group" => reply(&service.merge_group(args(body)?).await?),
242 "cancel" => reply(&service.cancel(args(body)?).await?),
243 "rerun" => reply(&service.rerun(args(body)?).await?),
244 "set_workflow_enabled" => reply(&service.set_workflow_enabled(args(body)?).await?),
245 "settings" => reply(&service.settings(args(body)?).await?),
246 "set_setting" => reply(&service.set_setting(args(body)?).await?),
247 "delete_setting" => reply(&service.delete_setting(args(body)?).await?),
248 "resolve_settings" => reply(&service.resolve_settings(args(body)?).await?),
249 "job_spec" => reply(&service.job_spec(args(body)?).await?),
250 "job_auth" => reply(&service.job_auth(args(body)?).await?),
251 "job_action" => reply(&service.job_action(args(body)?).await?),
252 "check_runs" => reply(&service.check_runs(args(body)?).await?),
253 "job_report" => reply(&service.job_report(args(body)?).await?),
254 // actions/cache, through the API with the job's token.
255 "cache_lookup" => reply(&service.cache_lookup(args(body)?).await?),
256 "cache_reserve" => reply(&service.cache_reserve(args(body)?).await?),
257 "cache_commit" => reply(&service.cache_commit(args(body)?).await?),
258 "cache_abort" => reply(&service.cache_abort(args(body)?).await?),
259 "cache_upload" => reply(&service.cache_upload(args(body)?).await?),
260 // Artifacts: a job's side, with its token or its runtime token.
261 "artifact_reserve" => reply(&service.artifact_reserve(args(body)?).await?),
262 "artifact_commit" => reply(&service.artifact_commit(args(body)?).await?),
263 "artifact_abort" => reply(&service.artifact_abort(args(body)?).await?),
264 "job_artifacts" => reply(&service.job_artifacts(args(body)?).await?),
265 "job_artifact" => reply(&service.job_artifact(args(body)?).await?),
266 "job_delete_artifact" => reply(&service.job_delete_artifact(args(body)?).await?),
267 // Artifacts: people's side.
268 "artifacts" => reply(&service.artifacts(args(body)?).await?),
269 "artifact" => reply(&service.artifact(args(body)?).await?),
270 "artifact_download" => reply(&service.artifact_download(args(body)?).await?),
271 "delete_artifact" => reply(&service.delete_artifact(args(body)?).await?),
272 "artifact_retention" => reply(&service.artifact_retention(args(body)?).await?),
273 // The toolkit's protocols: the runtime token, OIDC, blob links.
274 "runtime_auth" => reply(&service.runtime_auth(args(body)?).await?),
275 "oidc_claims" => reply(&service.oidc_claims(args(body)?).await?),
276 "blob_sign" => reply(&service.blob_sign(args(body)?).await?),
277 "blob_open" => reply(&service.blob_open(args(body)?).await?),
278 "blob_part" => reply(&service.blob_part(args(body)?).await?),
279 "blob_parts" => reply(&service.blob_parts(args(body)?).await?),
280 "blob_done" => reply(&service.blob_done(args(body)?).await?),
281 // Self-hosted runners: people's side.
282 "runners" => reply(&service.runners(args(body)?).await?),
283 "create_registration_token" => reply(&service.create_registration_token(args(body)?).await?),
284 "remove_runner" => reply(&service.remove_runner(args(body)?).await?),
285 "runner_groups" => reply(&service.runner_groups(args(body)?).await?),
286 "set_runner_group" => reply(&service.set_runner_group(args(body)?).await?),
287 "delete_runner_group" => reply(&service.delete_runner_group(args(body)?).await?),
288 "runner_settings" => reply(&service.runner_settings(args(body)?).await?),
289 "set_runner_settings" => reply(&service.set_runner_settings(args(body)?).await?),
290 // The runner's own side, through the API with its credential.
291 "runner_register" => reply(&service.runner_register(args(body)?).await?),
292 "runner_poll" => reply(&service.runner_poll(args(body)?).await?),
293 "runner_finished" => reply(&service.runner_finished(args(body)?).await?),
294 "runner_remove_self" => reply(&service.runner_remove_self(args(body)?).await?),
295 // Agent work, from the runner service.
296 "runner_route" => reply(&service.runner_route(args(body)?).await?),
297 "stuck_jobs" => reply(&service.stuck_jobs(args(body)?).await?),
298 "runner_activity" => reply(&service.runner_activity(args(body)?).await?),
299 "enqueue_task" => reply(&service.enqueue_task(args(body)?).await?),
300 "cancel_task" => reply(&service.cancel_task(args(body)?).await?),
301 _ => Response::error("Unknown method", 404),
302 }
303}
304
305/// Events from the bus, on this service's own queue.
306#[event(queue)]
307async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
308 let service = Actions::new(&env)?;
309 for message in batch.messages()? {
310 // A workspace renamed: its rows move to the slug it has now.
311 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), rename::STATEMENTS).await? {
312 message.ack();
313 continue;
314 }
315 // A repository renamed or transferred: its rows follow its new path.
316 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), rename::TRANSFERRED).await? {
317 message.ack();
318 continue;
319 }
320 // A workspace deleted: what it kept for itself goes.
321 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), rename::DELETED).await? {
322 message.ack();
323 continue;
324 }
325 // A repository purged: every row kept for it goes.
326 if g1t_kit::lifecycle::on_purged(&env.d1("DB")?, message.body(), rename::PURGED).await? {
327 message.ack();
328 continue;
329 }
330 // A repository deleted or archived: its runs stop.
331 if let Some(repo_id) = plan::stops_runs(message.body()) {
332 if let Err(error) = service.stop_runs(&repo_id).await {
333 worker::console_error!("actions: event {} failed: {error}", message.body().id);
334 message.retry();
335 continue;
336 }
337 message.ack();
338 continue;
339 }
340 if let Err(error) = service.on_event(message.body()).await {
341 worker::console_error!("actions: event {} failed: {error}", message.body().id);
342 message.retry();
343 continue;
344 }
345 message.ack();
346 }
347 Ok(())
348}
349
350/// Every minute: schedules that fire, jobs waiting for room, and jobs
351/// whose sandbox went quiet.
352#[event(scheduled)]
353async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
354 match Actions::new(&env) {
355 Ok(service) => {
356 if let Err(error) = service.on_minute(g1t_kit::now_ms()).await {
357 worker::console_error!("actions: the sweep failed: {error}");
358 }
359 }
360 Err(error) => worker::console_error!("actions: could not start: {error}"),
361 }
362}